field_ios

In-the-field app for Radroots on iOS
git clone https://radroots.dev/git/field_ios.git
Log | Files | Refs | README | LICENSE

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 }