TeraStoreObservation.swift (5035B)
1 import Foundation 2 3 /// Stable configuration inputs, excluding changing service observation evidence. 4 struct TeraPresentationConfiguration: Equatable { 5 let publicKey: String 6 let relayProfile: String? 7 let context: TeraLocalNetwork 8 let blossom: TeraBlossomConfigurationStatus? 9 10 init(snapshot: TeraRuntimeSnapshot) { 11 publicKey = snapshot.identity.publicKeyHex 12 relayProfile = snapshot.relay?.profile 13 context = .defaultContext(snapshot: snapshot) 14 blossom = snapshot.blossomConfiguration 15 } 16 } 17 18 /// At most one entry per fixed notification domain; a gap subsumes every domain. 19 struct TeraObservationBatch: Equatable { 20 private(set) var domains: Set<TeraRuntimeChangeKind> = [] 21 private(set) var resnapshot = false 22 23 static var currentState: Self { 24 Self(resnapshot: true) 25 } 26 27 var isEmpty: Bool { 28 !resnapshot && domains.isEmpty 29 } 30 31 func contains(anyOf kinds: Set<TeraRuntimeChangeKind>) -> Bool { 32 resnapshot || !domains.isDisjoint(with: kinds) 33 } 34 35 mutating func insert(_ change: TeraRuntimeChange) { 36 if change.delivery == .resnapshotRequired || change.revision.rawValue == nil { 37 self = .currentState 38 } else if !resnapshot, change.kind != .initial { 39 domains.insert(change.kind) 40 } 41 } 42 } 43 44 /// One stream and one refresh worker per store. Ingestion never waits for a 45 /// query; hints arriving during that query retain a bounded final refresh. 46 @MainActor 47 final class TeraStoreObservation { 48 private var generation = TeraSessionGeneration.initial 49 private var task: Task<Void, Never>? 50 private var refreshTask: Task<Void, Never>? 51 private var pending = TeraObservationBatch() 52 53 var isActive: Bool { 54 task != nil 55 } 56 57 deinit { 58 task?.cancel() 59 refreshTask?.cancel() 60 } 61 62 func stop() { 63 generation = generation.invalidated() 64 task?.cancel() 65 task = nil 66 refreshTask?.cancel() 67 refreshTask = nil 68 pending = TeraObservationBatch() 69 } 70 71 func start( 72 client: TeraRuntimeClient, 73 buffer: (capacity: Int, delay: @Sendable (UInt32) async throws -> Void), 74 state: @escaping @MainActor (TeraRuntimeObservationState) -> Void, 75 accepts: @escaping @MainActor (TeraRuntimeChange) -> Bool, 76 refresh: @escaping @MainActor (TeraObservationBatch) async -> Void 77 ) { 78 guard task == nil, generation.isActive else { return } 79 generation = generation.invalidated() 80 guard generation.isActive else { return } 81 let requested = generation 82 task = Task { [weak self] in 83 var attempt: UInt32 = 0 84 while self?.isCurrent(requested) == true { 85 state(.subscribing(attempt: attempt == .max ? .max : attempt + 1)) 86 var message = TeraUserMessages.text(.runtimeObservationUnavailable) 87 do { 88 let changes = try await client.changes(bufferCapacity: buffer.capacity) 89 guard self?.isCurrent(requested) == true else { break } 90 state(.active) 91 self?.requireCurrentState(requested, refresh: refresh) 92 for await value in changes { 93 guard self?.isCurrent(requested) == true else { break } 94 attempt = 0 95 if accepts(value) { 96 self?.pending.insert(value) 97 self?.startRefresh(requested, refresh: refresh) 98 } 99 } 100 } catch { 101 message = TeraUserMessages.text(for: error, fallback: .runtimeObservationUnavailable) 102 } 103 guard self?.isCurrent(requested) == true else { break } 104 attempt = attempt == .max ? .max : attempt + 1 105 state(.retrying(attempt: attempt, message: message)) 106 do { try await buffer.delay(attempt) } catch { break } 107 } 108 self?.finish(requested, state: state) 109 } 110 } 111 112 private func requireCurrentState( 113 _ requested: TeraSessionGeneration, 114 refresh: @escaping @MainActor (TeraObservationBatch) async -> Void 115 ) { 116 pending = .currentState 117 startRefresh(requested, refresh: refresh) 118 } 119 120 private func startRefresh( 121 _ requested: TeraSessionGeneration, 122 refresh: @escaping @MainActor (TeraObservationBatch) async -> Void 123 ) { 124 guard refreshTask == nil, !pending.isEmpty, isCurrent(requested) else { return } 125 refreshTask = Task { [weak self] in 126 while let batch = self?.takePending(requested) { 127 await refresh(batch) 128 } 129 guard self?.generation == requested else { return } 130 self?.refreshTask = nil 131 } 132 } 133 134 private func takePending(_ requested: TeraSessionGeneration) -> TeraObservationBatch? { 135 guard isCurrent(requested), !pending.isEmpty else { return nil } 136 let batch = pending 137 pending = TeraObservationBatch() 138 return batch 139 } 140 141 private func finish( 142 _ requested: TeraSessionGeneration, state: @MainActor (TeraRuntimeObservationState) -> Void 143 ) { 144 guard generation == requested else { return } 145 let cancelled = Task.isCancelled 146 stop() 147 if cancelled { 148 state(.stopped) 149 } 150 } 151 152 private func isCurrent(_ requested: TeraSessionGeneration) -> Bool { 153 generation == requested && generation.isActive && !Task.isCancelled 154 } 155 }