jsdw_ios/QuickLocation/Manager/MQTT/MQTTService.swift

1044 lines
38 KiB
Swift
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

//
// MQTTService.swift
// QuickLocation
//
// Created by 八条 on 2026/6/12.
//
import Foundation
import CocoaMQTT
import UIKit
import CoreLocation
import Network
import RxSwift
// MARK: - MQTT 模型
/// 位置上报类型
enum MqttType: String, Codable {
case track = "track" // 追踪(获取)和更新用户状态
case disconnect = "disconnect"
case message = "message"
case join = "join" // 加入群
case leave = "leave" // 离开群
case dismiss = "dismiss" // 解散群
case kick = "kick" // 踢人
case signIn = "signIn" // 签到
case sos = "sos" // 求助
case sosQuit = "sosquit" // 取消求助
case bubbleQuit = "bubblequit" // 退出气泡
case needTrack = "needtrack" // 追踪上报
case emote = "emote" // 接收表情
case lockApp = "lockApp" // 远程锁 App
case appUnlock = "appUnlock" // 远程解锁 App
case lockAppsChanged = "lockAppsChanged" // App 锁状态已同步
}
/// 单点位置
struct Points: Codable {
let lat: Double
let lon: Double
let addr: String
let time: Int64 // 毫秒时间戳
let speed: Double
let bearing: Double
let altitude: Double
let accuracy: Double
}
/// 位置上报数据(发送)
struct MqttLocation: Codable {
let battery: String
let points: [Points]
}
/// 签到上报数据
struct MqttSignIn: Codable {
let topic: String?
let latitude: Double
let longitude: Double
let address: String
}
/// MQTT 收到的位置消息(接收解析用)
struct MqttIncomingMessage: Decodable {
let type: String?
let data: MqttIncomingData?
let extra: String?
}
struct MqttIncomingData: Decodable {
let points: [Points]?
let battery: String?
let index: Int?
let group_key: String?
let user_id: String?
let lock_app: MqttLockAppBody?
let app_unlock: MqttLockAppBody?
let sender: String?
}
struct MqttLockAppBody: Decodable {
let os: String?
let group_key: String?
let group_name: String?
let user_id: String?
let icon_index: Int?
let message: String?
let lock_time: Int64?
let app_icon: String?
let token: MqttFlexibleStringArray?
}
private struct MqttLockAppsChangedData: Codable, Equatable {
let userId: String
let groupKey: String
enum CodingKeys: String, CodingKey {
case userId = "user_id"
case groupKey = "group_key"
}
}
private struct MqttLockAppsChangedPayload: Codable {
let type: String
let data: MqttLockAppsChangedData
let extra: String
}
enum MqttFlexibleStringArray: Decodable {
case one(String)
case many([String])
var values: [String] {
switch self {
case .one(let value):
return value.isEmpty ? [] : [value]
case .many(let values):
return values.filter { !$0.isEmpty }
}
}
init(from decoder: Decoder) throws {
let container = try decoder.singleValueContainer()
if let values = try? container.decode([String].self) {
self = .many(values)
return
}
if let value = try? container.decode(String.self) {
self = .one(value)
return
}
if let value = try? container.decode(Int.self) {
self = .one(String(value))
return
}
self = .many([])
}
}
extension PhoneLockRecord {
init(mqtt body: MqttLockAppBody) {
self.init()
groupKey = body.group_key ?? ""
groupName = body.group_name ?? ""
userId = body.user_id ?? ""
tokens = body.token?.values ?? []
os = body.os ?? ""
iconIndex = body.icon_index ?? 0
message = body.message ?? ""
lockTime = body.lock_time ?? 0
appIcon = body.app_icon ?? ""
}
}
private struct LegacyMqttAppUnlockReceiptBody: Codable {
let os: String
let group_key: String
let group_name: String
let user_id: String
let token: [String]
}
private struct LegacyMqttAppUnlockReceiptData: Codable {
let app_unlock: LegacyMqttAppUnlockReceiptBody
}
private struct LegacyMqttAppUnlockReceiptPayload: Codable {
let type: String
let data: LegacyMqttAppUnlockReceiptData
let extra: String
}
private struct LegacyPendingAppUnlockReceipt: Codable {
let identifier: String
let userId: String
let payload: String
private enum CodingKeys: String, CodingKey {
case identifier
case userId
case payload
}
init(from decoder: Decoder) throws {
let container = try decoder.container(keyedBy: CodingKeys.self)
identifier = try container.decodeIfPresent(String.self, forKey: .identifier) ?? UUID().uuidString
userId = try container.decode(String.self, forKey: .userId)
payload = try container.decode(String.self, forKey: .payload)
}
}
private struct PendingAppUnlockConfirmation: Codable, Equatable {
let identifier: String
let userId: String
let os: String
let groupKey: String
let tokens: [String]
init(
identifier: String = UUID().uuidString,
userId: String,
os: String,
groupKey: String,
tokens: [String]
) {
self.identifier = identifier
self.userId = userId
self.os = os
self.groupKey = groupKey
self.tokens = tokens
}
}
final class AppUnlockCoordinator {
static let shared = AppUnlockCoordinator()
private let pendingConfirmationsKey = "AppUnlockPendingConfirmations"
private let legacyPendingReceiptsKey = "AppUnlockPendingReceipts"
private var pendingConfirmations: [PendingAppUnlockConfirmation]
private var inFlightConfirmation: PendingAppUnlockConfirmation?
private var allowRequestUserId = ""
private var completedAllowUserIds: Set<String> = []
private var observedUserId = ""
private var lockStateRevision: UInt64 = 0
private var locallyUnlockedTokensByUserId: [String: Set<String>] = [:]
private var locallyUnlockedAllUserIds: Set<String> = []
private var localUnlockSuppressionExpirations: [String: Date] = [:]
private var allowDisposable: Disposable?
private var confirmationDisposable: Disposable?
private let confirmationNetworkMonitor = NWPathMonitor()
private let confirmationNetworkQueue = DispatchQueue(label: "com.quicklocation.app-unlock-confirmation")
private var observers: [NSObjectProtocol] = []
private var isStarted = false
private init() {
if let data = UserDefaults.standard.data(forKey: pendingConfirmationsKey),
let confirmations = try? JSONDecoder().decode([PendingAppUnlockConfirmation].self, from: data) {
pendingConfirmations = confirmations
} else if let data = UserDefaults.standard.data(forKey: legacyPendingReceiptsKey),
let receipts = try? JSONDecoder().decode([LegacyPendingAppUnlockReceipt].self, from: data) {
pendingConfirmations = receipts.compactMap(Self.makeConfirmation(from:))
if let migratedData = try? JSONEncoder().encode(pendingConfirmations) {
UserDefaults.standard.set(migratedData, forKey: pendingConfirmationsKey)
UserDefaults.standard.removeObject(forKey: legacyPendingReceiptsKey)
}
} else {
pendingConfirmations = []
}
for confirmation in pendingConfirmations {
rememberLocalUnlock(tokens: confirmation.tokens, userId: confirmation.userId)
}
}
func start() {
performOnMain { [weak self] in
guard let self, !self.isStarted else { return }
self.isStarted = true
self.confirmationNetworkMonitor.pathUpdateHandler = { [weak self] path in
guard path.status == .satisfied else { return }
self?.performOnMain { [weak self] in
self?.retryPendingConfirmations()
}
}
self.confirmationNetworkMonitor.start(queue: self.confirmationNetworkQueue)
self.observers.append(
NotificationCenter.default.addObserver(
forName: UIApplication.didBecomeActiveNotification,
object: nil,
queue: .main
) { [weak self] _ in
self?.checkOfflineUnlockForCurrentUser()
self?.retryPendingConfirmations()
}
)
self.observers.append(
NotificationCenter.default.addObserver(
forName: UIApplication.didEnterBackgroundNotification,
object: nil,
queue: .main
) { [weak self] _ in
let userId = AppContextManager.shared.userId.trimmed
self?.completedAllowUserIds.remove(userId)
}
)
self.observers.append(
NotificationCenter.default.addObserver(
forName: .RefreshUserConfigNotification,
object: nil,
queue: .main
) { [weak self] _ in
self?.checkOfflineUnlockForCurrentUser()
self?.retryPendingConfirmations()
}
)
self.checkOfflineUnlockForCurrentUser()
self.retryPendingConfirmations()
}
}
func checkOfflineUnlockForCurrentUser() {
performOnMain { [weak self] in
guard let self else { return }
let userId = AppContextManager.shared.userId.trimmed
guard !userId.isEmpty else { return }
if self.observedUserId != userId {
self.cancelConfirmationForDifferentUser(userId)
self.observedUserId = userId
self.completedAllowUserIds.remove(userId)
self.retryPendingConfirmations()
}
guard !self.completedAllowUserIds.contains(userId) else { return }
if !self.allowRequestUserId.isEmpty, self.allowRequestUserId != userId {
self.allowDisposable?.dispose()
self.allowDisposable = nil
self.allowRequestUserId = ""
}
guard self.allowRequestUserId.isEmpty else { return }
self.allowRequestUserId = userId
self.allowDisposable = UserService.phoneUnlockAllow()
.observe(on: MainScheduler.instance)
.subscribe(onNext: { [weak self] response in
guard let self else { return }
self.finishAllowRequest(for: userId)
guard AppContextManager.shared.userId.trimmed == userId else {
self.checkOfflineUnlockForCurrentUser()
return
}
guard response.code == "0" else {
print("[AppUnlock] unlock allow failed: \(response.message ?? "unknown error")")
return
}
self.completedAllowUserIds.insert(userId)
for record in response.model?.locks ?? [] {
self.applyUnlockAndEnqueueConfirmation(record: record, userId: userId)
}
}, onError: { [weak self] error in
self?.finishAllowRequest(for: userId)
print("[AppUnlock] unlock allow request failed: \(error.gatewayMessage ?? error.localizedDescription)")
})
}
}
func accountDidLogout(userId: String) {
performOnMain { [weak self] in
guard let self else { return }
let userId = userId.trimmed
self.completedAllowUserIds.remove(userId)
if self.allowRequestUserId == userId {
self.allowDisposable?.dispose()
self.allowDisposable = nil
self.allowRequestUserId = ""
}
if self.observedUserId == userId {
self.observedUserId = ""
}
if self.inFlightConfirmation?.userId == userId {
self.confirmationDisposable?.dispose()
self.confirmationDisposable = nil
self.inFlightConfirmation = nil
}
self.locallyUnlockedTokensByUserId.removeValue(forKey: userId)
self.locallyUnlockedAllUserIds.remove(userId)
self.localUnlockSuppressionExpirations.removeValue(forKey: userId)
}
}
@discardableResult
func handleIncoming(topic: String, payload: String?) -> Bool {
guard let payload,
let data = payload.data(using: .utf8),
let message = try? JSONDecoder().decode(MqttIncomingMessage.self, from: data),
message.type == MqttType.appUnlock.rawValue else { return false }
print("📩 收到消息 -> 主题:\(topic),内容:\(payload)")
performOnMain { [weak self] in
guard let self else { return }
guard let body = message.data?.app_unlock else { return }
let currentUserId = AppContextManager.shared.userId.trimmed
let topicUserId = topic.replacingOccurrences(of: "smartdrive/", with: "").trimmed
let targetUserId = body.user_id?.trimmed ?? ""
guard !currentUserId.isEmpty,
targetUserId == currentUserId || (targetUserId.isEmpty && topicUserId == currentUserId) else { return }
var record = PhoneLockRecord(mqtt: body)
record.userId = currentUserId
self.applyUnlockAndEnqueueConfirmation(record: record, userId: currentUserId)
}
return true
}
func applyUnlock(tokens: [String]) {
guard Thread.isMainThread else {
DispatchQueue.main.async { [weak self] in
self?.applyUnlock(tokens: tokens)
}
return
}
lockStateRevision &+= 1
rememberLocalUnlock(tokens: tokens, userId: AppContextManager.shared.userId.trimmed)
if #available(iOS 16.0, *) {
AppRestrictManager.shared.applyRemoteUnlock(tokens: tokens)
}
if tokens.isEmpty {
PhoneLockSession.currentLocks = []
} else {
let tokenSet = Set(tokens)
PhoneLockSession.currentLocks = PhoneLockSession.currentLocks.compactMap { record in
var updated = record
updated.tokens = record.tokens.filter { !tokenSet.contains($0) }
return updated.tokens.isEmpty ? nil : updated
}
}
if PhoneLockSession.currentLocks.isEmpty {
LockedAppPopView.dismiss()
}
}
func registerIncomingLock(tokens: [String]) {
guard Thread.isMainThread else {
DispatchQueue.main.async { [weak self] in
self?.registerIncomingLock(tokens: tokens)
}
return
}
let userId = AppContextManager.shared.userId.trimmed
guard !userId.isEmpty else { return }
lockStateRevision &+= 1
locallyUnlockedAllUserIds.remove(userId)
guard var unlockedTokens = locallyUnlockedTokensByUserId[userId] else { return }
unlockedTokens.subtract(tokens)
if unlockedTokens.isEmpty {
locallyUnlockedTokensByUserId.removeValue(forKey: userId)
} else {
locallyUnlockedTokensByUserId[userId] = unlockedTokens
}
if locallyUnlockedTokensByUserId[userId] == nil {
localUnlockSuppressionExpirations.removeValue(forKey: userId)
}
}
func currentLockStateRevision() -> UInt64 {
if !Thread.isMainThread {
assertionFailure("Lock state revision must be read on the main thread")
}
return lockStateRevision
}
func canApplyLockResponse(startedAt revision: UInt64, userId: String) -> Bool {
if !Thread.isMainThread {
assertionFailure("Lock response must be checked on the main thread")
}
return revision == lockStateRevision
&& userId == AppContextManager.shared.userId.trimmed
}
func locksForRestoration(
_ locks: [PhoneLockRecord],
startedAt revision: UInt64,
userId: String
) -> [PhoneLockRecord]? {
guard canApplyLockResponse(startedAt: revision, userId: userId) else { return nil }
removeExpiredLocalUnlockSuppression(for: userId)
let pendingScope = pendingUnlockScope(for: userId)
if locallyUnlockedAllUserIds.contains(userId) || pendingScope.unlockAll {
return []
}
let unlockedTokens = locallyUnlockedTokensByUserId[userId, default: []]
.union(pendingScope.tokens)
guard !unlockedTokens.isEmpty else { return locks }
return locks.compactMap { record in
var filteredRecord = record
filteredRecord.tokens = record.tokens.filter { !unlockedTokens.contains($0) }
return filteredRecord.tokens.isEmpty ? nil : filteredRecord
}
}
private func applyUnlockAndEnqueueConfirmation(record: PhoneLockRecord, userId: String) {
applyUnlock(tokens: record.tokens)
let confirmation = PendingAppUnlockConfirmation(
userId: userId,
os: record.os.trimmed.isEmpty ? "ios" : record.os.trimmed,
groupKey: record.groupKey,
tokens: record.tokens
)
let isPending = pendingConfirmations.contains {
$0.userId == confirmation.userId
&& $0.os == confirmation.os
&& $0.groupKey == confirmation.groupKey
&& $0.tokens == confirmation.tokens
}
if !isPending {
pendingConfirmations.append(confirmation)
savePendingConfirmations()
}
retryPendingConfirmations()
}
private func retryPendingConfirmations() {
guard Thread.isMainThread else {
DispatchQueue.main.async { [weak self] in
self?.retryPendingConfirmations()
}
return
}
let userId = AppContextManager.shared.userId.trimmed
guard !userId.isEmpty else { return }
cancelConfirmationForDifferentUser(userId)
guard inFlightConfirmation == nil,
let confirmation = pendingConfirmations.first(where: { $0.userId == userId }) else { return }
inFlightConfirmation = confirmation
confirmationDisposable = UserService.phoneUnlockAllowConfirm(
os: confirmation.os,
groupKey: confirmation.groupKey,
userId: confirmation.userId,
tokens: confirmation.tokens
)
.observe(on: MainScheduler.instance)
.subscribe(onNext: { [weak self] response in
guard let self,
self.inFlightConfirmation?.identifier == confirmation.identifier else { return }
self.inFlightConfirmation = nil
self.confirmationDisposable = nil
guard response.code == "0" else {
print("[AppUnlock] unlock confirmation failed: \(response.message ?? "unknown error")")
return
}
MQTTService.shared.reportLockAppsChanged(
userId: confirmation.userId,
groupKey: confirmation.groupKey
)
self.pendingConfirmations.removeAll { $0.identifier == confirmation.identifier }
self.savePendingConfirmations()
self.retryPendingConfirmations()
}, onError: { [weak self] error in
guard let self,
self.inFlightConfirmation?.identifier == confirmation.identifier else { return }
self.inFlightConfirmation = nil
self.confirmationDisposable = nil
print("[AppUnlock] unlock confirmation request failed: \(error.gatewayMessage ?? error.localizedDescription)")
})
}
private func finishAllowRequest(for userId: String) {
guard allowRequestUserId == userId else { return }
allowRequestUserId = ""
allowDisposable = nil
}
private func cancelConfirmationForDifferentUser(_ userId: String) {
guard let confirmation = inFlightConfirmation,
confirmation.userId != userId else { return }
confirmationDisposable?.dispose()
confirmationDisposable = nil
inFlightConfirmation = nil
}
private func savePendingConfirmations() {
guard let data = try? JSONEncoder().encode(pendingConfirmations) else { return }
UserDefaults.standard.set(data, forKey: pendingConfirmationsKey)
}
private func rememberLocalUnlock(tokens: [String], userId: String) {
guard !userId.isEmpty else { return }
localUnlockSuppressionExpirations[userId] = Date().addingTimeInterval(120)
if tokens.isEmpty {
locallyUnlockedAllUserIds.insert(userId)
locallyUnlockedTokensByUserId.removeValue(forKey: userId)
return
}
locallyUnlockedTokensByUserId[userId, default: []].formUnion(tokens)
}
private func removeExpiredLocalUnlockSuppression(for userId: String) {
guard let expiresAt = localUnlockSuppressionExpirations[userId],
expiresAt <= Date() else { return }
locallyUnlockedTokensByUserId.removeValue(forKey: userId)
locallyUnlockedAllUserIds.remove(userId)
localUnlockSuppressionExpirations.removeValue(forKey: userId)
}
private func pendingUnlockScope(for userId: String) -> (unlockAll: Bool, tokens: Set<String>) {
var unlockAll = false
var tokens: Set<String> = []
for confirmation in pendingConfirmations where confirmation.userId == userId {
if confirmation.tokens.isEmpty {
unlockAll = true
} else {
tokens.formUnion(confirmation.tokens)
}
}
return (unlockAll, tokens)
}
private static func makeConfirmation(
from receipt: LegacyPendingAppUnlockReceipt
) -> PendingAppUnlockConfirmation? {
guard let data = receipt.payload.data(using: .utf8),
let payload = try? JSONDecoder().decode(LegacyMqttAppUnlockReceiptPayload.self, from: data) else {
return nil
}
let body = payload.data.app_unlock
let userId = receipt.userId.trimmed.isEmpty ? body.user_id.trimmed : receipt.userId.trimmed
guard !userId.isEmpty else { return nil }
return PendingAppUnlockConfirmation(
identifier: receipt.identifier,
userId: userId,
os: body.os.trimmed.isEmpty ? "ios" : body.os.trimmed,
groupKey: body.group_key,
tokens: body.token
)
}
private func performOnMain(_ work: @escaping () -> Void) {
if Thread.isMainThread {
work()
} else {
DispatchQueue.main.async(execute: work)
}
}
}
// MARK: - MQTTService
/// MQTT 5.0 服务,管理连接、订阅和消息收发
final class MQTTService: NSObject {
static let shared = MQTTService()
private var mqtt: CocoaMQTT5?
private(set) var isConnected = false
// MARK: - 连接状态回调
var onConnected: (() -> Void)?
var onDisconnected: (() -> Void)?
/// 全局接收回调(所有未匹配 topicCallback 的消息)
var onMessageReceived: ((CocoaMQTT5Message, UInt16, MqttPublishProperties) -> Void)?
/// 按 topic 订阅的回调(优先匹配)
private var topicCallbacks: [String: (CocoaMQTT5Message) -> Void] = [:]
private var lockStateSubscriptionsByOwner: [UUID: Set<String>] = [:]
private var pendingLockAppsChangedEvents: [MqttLockAppsChangedData] = []
// MARK: - 配置
private var host: String { "emqx.batiao8.com" }
private var port: UInt16 { 1883 }
/// 当前 clientID,切换用户时可更新后重连
private(set) var clientID: String = ""
private var userName = "batiao"
private var password = "Batiao12B"
private var topic = "smartdrive/"
private var lockStateSubscriptionTopics: Set<String> {
lockStateSubscriptionsByOwner.values.reduce(into: Set<String>()) { result, topics in
result.formUnion(topics)
}
}
private func postLockDistractAppsChangeIfNeeded(topic: String, payload: String?) {
guard let payload,
let data = payload.data(using: .utf8),
let message = try? JSONDecoder().decode(MqttIncomingMessage.self, from: data),
let type = message.type else { return }
let topicUserId = topic.replacingOccurrences(of: self.topic, with: "").trimmed
let userId: String
let groupKey: String
switch type {
case MqttType.lockApp.rawValue:
userId = message.data?.lock_app?.user_id?.trimmed ?? topicUserId
groupKey = message.data?.lock_app?.group_key?.trimmed ?? ""
case MqttType.appUnlock.rawValue:
userId = message.data?.app_unlock?.user_id?.trimmed ?? topicUserId
groupKey = message.data?.app_unlock?.group_key?.trimmed ?? ""
case MqttType.lockAppsChanged.rawValue:
userId = message.data?.user_id?.trimmed ?? topicUserId
groupKey = message.data?.group_key?.trimmed ?? ""
default:
return
}
postLockDistractAppsChange(userId: userId, groupKey: groupKey)
}
private func postLockDistractAppsChange(userId: String, groupKey: String) {
let post = {
NotificationCenter.default.post(
name: .lockDistractAppsDidChange,
object: nil,
userInfo: [
LockDistractAppsChangeUserInfoKey.userId: userId,
LockDistractAppsChangeUserInfoKey.groupKey: groupKey
]
)
}
if Thread.isMainThread {
post()
} else {
DispatchQueue.main.async(execute: post)
}
}
private func postGroupDataChangeIfNeeded(payload: String?) {
guard let payload,
let data = payload.data(using: .utf8),
let message = try? JSONDecoder().decode(MqttIncomingMessage.self, from: data),
let operation = message.type,
["join", "leave", "dismiss", "kick"].contains(operation) else { return }
var userInfo: [AnyHashable: Any] = [
GroupDataChangeUserInfoKey.operation: operation,
GroupDataChangeUserInfoKey.source: GroupDataChangeSource.mqtt
]
if let groupKey = message.data?.group_key, !groupKey.isEmpty {
userInfo[GroupDataChangeUserInfoKey.groupKey] = groupKey
}
if let userId = message.data?.user_id, !userId.isEmpty {
userInfo[GroupDataChangeUserInfoKey.userId] = userId
}
let post = {
NotificationCenter.default.post(
name: .RefreshGroupInfoNotification,
object: nil,
userInfo: userInfo
)
}
if Thread.isMainThread {
post()
} else {
DispatchQueue.main.sync(execute: post)
}
}
override private init() {
super.init()
UIDevice.current.isBatteryMonitoringEnabled = true
}
// MARK: - 连接
func connect() {
guard !isConnected else { return }
if clientID.isEmpty {
clientID = "smartdrive_\(AppContextManager.shared.userId)"
}
let mqtt = CocoaMQTT5(clientID: clientID, host: host, port: port)
mqtt.username = userName
mqtt.password = password
mqtt.cleanSession = true
mqtt.keepAlive = 30
mqtt.autoReconnect = true
mqtt.autoReconnectTimeInterval = 5
mqtt.delegate = self
mqtt.logLevel = .warning
// MQTT 5 连接属性(对应 Android 端配置)
let connProps = MqttConnectProperties()
connProps.sessionExpiryInterval = 30
connProps.receiveMaximum = 10
connProps.maximumPacketSize = 10240
connProps.topicAliasMaximum = 0
mqtt.connectProperties = connProps
// Will 消息
let will = CocoaMQTT5Message(topic: "willtopic", payload: [])
will.qos = .qos1
will.retained = false
mqtt.willMessage = will
self.mqtt = mqtt
_ = mqtt.connect()
}
// MARK: - 断开
func disconnect() {
mqtt?.disconnect()
isConnected = false
}
// MARK: - 切换用户
/// 更新 clientID 并重连(切换用户后调用)
func updateClientID(_ newID: String) {
clientID = newID
disconnect()
connect()
}
// MARK: - 订阅主题
/// - Parameters:
/// - topic: 主题
/// - qos: 服务质量
/// - callback: 可选,该 topic 的专用回调,收到消息时优先于此回调
func subscribe(topic: String, qos: CocoaMQTTQoS = .qos1, callback: ((CocoaMQTT5Message) -> Void)? = nil) {
let subscription = MqttSubscription(topic: topic, qos: qos)
subscription.noLocal = false
subscription.retainAsPublished = true
subscription.retainHandling = .none
mqtt?.subscribe([subscription])
if let cb = callback {
topicCallbacks[topic] = cb
}
}
func updateLockStateSubscriptions(owner: UUID, userIds: [String]) {
performOnMain { [weak self] in
guard let self else { return }
let previousTopics = self.lockStateSubscriptionTopics
let nextOwnerTopics = Set(
userIds
.map(\.trimmed)
.filter { !$0.isEmpty }
.map { "\(self.topic)\($0)" }
)
if nextOwnerTopics.isEmpty {
self.lockStateSubscriptionsByOwner.removeValue(forKey: owner)
} else {
self.lockStateSubscriptionsByOwner[owner] = nextOwnerTopics
}
let nextTopics = self.lockStateSubscriptionTopics
for subscriptionTopic in nextTopics.subtracting(previousTopics) {
self.subscribe(topic: subscriptionTopic)
}
for subscriptionTopic in previousTopics.subtracting(nextTopics)
where self.topicCallbacks[subscriptionTopic] == nil
&& !self.isRequiredBaseSubscription(subscriptionTopic) {
self.mqtt?.unsubscribe(subscriptionTopic)
}
}
}
// MARK: - 取消订阅
func unsubscribe(topic: String) {
mqtt?.unsubscribe(topic)
topicCallbacks.removeValue(forKey: topic)
}
/// 订阅指定圈子的所有成员位置 topic: smartdrive/<memberId>
func subscribeGroupMembers(_ memberIds: [String]) {
for id in memberIds {
subscribe(topic: "\(topic)\(id)")
}
}
/// 取消订阅上一批成员
func unsubscribeGroupMembers(_ memberIds: [String]) {
for id in memberIds {
let subscriptionTopic = "\(topic)\(id)"
topicCallbacks.removeValue(forKey: subscriptionTopic)
if !lockStateSubscriptionTopics.contains(subscriptionTopic),
!isRequiredBaseSubscription(subscriptionTopic) {
mqtt?.unsubscribe(subscriptionTopic)
}
}
}
func reportLockAppsChanged(userId: String, groupKey: String) {
performOnMain { [weak self] in
guard let self else { return }
let event = MqttLockAppsChangedData(
userId: userId.trimmed,
groupKey: groupKey.trimmed
)
guard !event.userId.isEmpty else { return }
self.postLockDistractAppsChange(userId: event.userId, groupKey: event.groupKey)
if !self.pendingLockAppsChangedEvents.contains(event) {
self.pendingLockAppsChangedEvents.append(event)
}
self.flushPendingLockAppsChangedEvents()
}
}
private func flushPendingLockAppsChangedEvents() {
guard isConnected, !pendingLockAppsChangedEvents.isEmpty else { return }
var unsentEvents: [MqttLockAppsChangedData] = []
for event in pendingLockAppsChangedEvents {
let payload = MqttLockAppsChangedPayload(
type: MqttType.lockAppsChanged.rawValue,
data: event,
extra: ""
)
guard let data = try? JSONEncoder().encode(payload),
let message = String(data: data, encoding: .utf8),
publish(topic: "\(topic)\(event.userId)", message: message) >= 0 else {
unsentEvents.append(event)
continue
}
}
pendingLockAppsChangedEvents = unsentEvents
}
private func isRequiredBaseSubscription(_ subscriptionTopic: String) -> Bool {
if subscriptionTopic == topic { return true }
let currentUserId = AppContextManager.shared.userId.trimmed
return !currentUserId.isEmpty && subscriptionTopic == "\(topic)\(currentUserId)"
}
// MARK: - 发布消息
@discardableResult
func publish(topic: String, message: String, qos: CocoaMQTTQoS = .qos1) -> Int {
let properties = MqttPublishProperties()
return mqtt?.publish(topic, withString: message, qos: qos, DUP: false, retained: false, properties: properties) ?? -1
}
@discardableResult
func publish(topic: String, data: Data, qos: CocoaMQTTQoS = .qos1) -> Int {
let properties = MqttPublishProperties()
let message = CocoaMQTT5Message(topic: topic, payload: [UInt8](data))
return mqtt?.publish(message, DUP: false, retained: false, properties: properties) ?? -1
}
// MARK: - 位置上报
/// 构建并上报位置数据(格式与 Android 一致)
func reportLocation(lat: Double, lon: Double, addr: String,
speed: CLLocationSpeed, bearing: CLLocationDirection,
altitude: CLLocationDistance, accuracy: CLLocationAccuracy) {
let battery = UIDevice.batteryPercent.map(String.init) ?? ""
let point = Points(
lat: lat, lon: lon, addr: addr,
time: Int64(Date().timeIntervalSince1970 * 1000),
speed: speed, bearing: bearing,
altitude: altitude, accuracy: accuracy
)
let location = MqttLocation(battery: battery, points: [point])
guard let jsonData = try? JSONEncoder().encode(location),
let dataDict = try? JSONSerialization.jsonObject(with: jsonData) as? [String: Any]
else { return }
// 外层包装,data 是 JSON 对象而非字符串
let payload: [String: Any] = [
"type": MqttType.track.rawValue,
"extra": "",
"data": dataDict
]
publish(topic: "\(topic)\(AppContextManager.shared.userId)", message: payload.toJsonString())
}
// MARK: - 签到
func reportSignIn(lat: Double, lon: Double, addr: String) {
let signIn = MqttSignIn(topic: nil, latitude: lat, longitude: lon, address: addr)
guard let jsonData = try? JSONEncoder().encode(signIn),
let dataDict = try? JSONSerialization.jsonObject(with: jsonData) as? [String: Any]
else { return }
let payload: [String: Any] = [
"type": MqttType.signIn.rawValue,
"data": dataDict,
"extra": "",
"sort": ""
]
publish(topic: "\(topic)\(AppContextManager.shared.userId)", message: payload.toJsonString())
}
private func performOnMain(_ work: @escaping () -> Void) {
if Thread.isMainThread {
work()
} else {
DispatchQueue.main.async(execute: work)
}
}
}
// MARK: - CocoaMQTT5Delegate
extension MQTTService: CocoaMQTT5Delegate {
func mqtt5(_ mqtt5: CocoaMQTT5, didConnectAck ack: CocoaMQTTCONNACKReasonCode, connAckData: MqttDecodeConnAck?) {
isConnected = true
print("MQTT5 connected: \(ack)")
// 订阅基础 topic,接收 signIn/join/leave 等非位置消息
subscribe(topic: topic)
let userId = AppContextManager.shared.userId.trimmed
if !userId.isEmpty {
subscribe(topic: "\(topic)\(userId)")
}
for subscriptionTopic in Set(topicCallbacks.keys).union(lockStateSubscriptionTopics) {
subscribe(topic: subscriptionTopic)
}
onConnected?()
flushPendingLockAppsChangedEvents()
}
func mqtt5(_ mqtt5: CocoaMQTT5, didPublishMessage message: CocoaMQTT5Message, id: UInt16) {
print("MQTT5 published: \(message.topic)")
}
func mqtt5(_ mqtt5: CocoaMQTT5, didPublishAck id: UInt16, pubAckData: MqttDecodePubAck?) {
print("MQTT5 publish ack: \(id)")
}
func mqtt5(_ mqtt5: CocoaMQTT5, didPublishRec id: UInt16, pubRecData: MqttDecodePubRec?) {}
func mqtt5(_ mqtt5: CocoaMQTT5, didReceiveMessage message: CocoaMQTT5Message, id: UInt16, publishData: MqttDecodePublish?) {
postGroupDataChangeIfNeeded(payload: message.string)
postLockDistractAppsChangeIfNeeded(topic: message.topic, payload: message.string)
if AppUnlockCoordinator.shared.handleIncoming(topic: message.topic, payload: message.string) {
return
}
// 优先 topic 专用回调
if let cb = topicCallbacks[message.topic] {
cb(message)
return
}
// 没有专用回调时走全局回调
if let payload = message.string {
print("MQTT5 received on \(message.topic): \(payload)")
}
onMessageReceived?(message, id, MqttPublishProperties())
}
func mqtt5(_ mqtt5: CocoaMQTT5, didSubscribeTopics success: NSDictionary, failed: [String], subAckData: MqttDecodeSubAck?) {
print("MQTT5 subscribe success: \(success), failed: \(failed)")
}
func mqtt5(_ mqtt5: CocoaMQTT5, didUnsubscribeTopics topics: [String], unsubAckData: MqttDecodeUnsubAck?) {
print("MQTT5 unsubscribe: \(topics)")
}
func mqtt5(_ mqtt5: CocoaMQTT5, didReceiveDisconnectReasonCode reasonCode: CocoaMQTTDISCONNECTReasonCode) {}
func mqtt5(_ mqtt5: CocoaMQTT5, didReceiveAuthReasonCode reasonCode: CocoaMQTTAUTHReasonCode) {}
func mqtt5DidPing(_ mqtt5: CocoaMQTT5) {}
func mqtt5DidReceivePong(_ mqtt5: CocoaMQTT5) {}
func mqtt5DidDisconnect(_ mqtt5: CocoaMQTT5, withError err: Error?) {
isConnected = false
print("MQTT5 disconnected: \(err?.localizedDescription ?? "")")
onDisconnected?()
}
}