TeraRuntimeInvalidation.swift (3423B)
1 import Foundation 2 3 enum TeraRuntimeChangeKind: Sendable, Equatable, Hashable { 4 case initial 5 case identity 6 case settings 7 case profile 8 case today 9 case drafts 10 case relay 11 case media 12 case lifecycle 13 } 14 15 enum TeraRuntimeChangeDelivery: Sendable, Equatable { 16 case change 17 case resnapshotRequired 18 } 19 20 struct TeraRuntimeChangeScope: Sendable, Equatable, Hashable { 21 let publicKey: String 22 let sourceGeneration: String 23 let context: TeraLocalNetwork? 24 } 25 26 /// A hint to query state, never a durable mutation receipt or host session ID. 27 /// Revisions are comparable within one epoch, domain and exact scope only. 28 struct TeraRuntimeChange: Sendable, Equatable { 29 let schemaVersion: UInt16 30 let scope: TeraRuntimeChangeScope 31 let epoch: String 32 let revision: TeraProjectionRevision 33 let delivery: TeraRuntimeChangeDelivery 34 let kind: TeraRuntimeChangeKind 35 let entityID: String? 36 37 func matches(_ configuration: TeraRuntimeLaunchConfiguration?) -> Bool { 38 guard let configuration else { return false } 39 return schemaVersion == 3 40 && Self.isHex(epoch, count: 32) 41 && Self.isHex(scope.publicKey, count: 64) 42 && Self.isHex(scope.sourceGeneration, count: 64) 43 && scope.publicKey == configuration.publicKeyHex 44 && scope.sourceGeneration == configuration.sourceGenerationHex 45 && (scope.context == nil || scope.context?.schemaVersion == 1) 46 } 47 48 func matches(context: TeraLocalNetwork?) -> Bool { 49 scope.context == nil || scope.context == context 50 } 51 52 func requiringResnapshot() -> Self { 53 Self( 54 schemaVersion: schemaVersion, 55 scope: TeraRuntimeChangeScope(publicKey: scope.publicKey, sourceGeneration: scope.sourceGeneration, context: nil), 56 epoch: epoch, revision: TeraProjectionRevision(rawValue: 0), 57 delivery: .resnapshotRequired, kind: .initial, entityID: nil 58 ) 59 } 60 61 /// Loss leaves an account-wide gap queued or already delivered, even if this 62 /// is the producer's final event. The existing buffer capacity stays fixed. 63 func yield(to continuation: AsyncStream<Self>.Continuation) { 64 if case .dropped = continuation.yield(self) { 65 continuation.yield(requiringResnapshot()) 66 } 67 } 68 69 private static func isHex(_ value: String, count: Int) -> Bool { 70 value.utf8.count == count && value.utf8.contains(where: { $0 != 48 }) 71 && value.utf8.allSatisfy { (48 ... 57).contains($0) || (97 ... 102).contains($0) } 72 } 73 } 74 75 /// One subscription retains at most one watermark for each of the nine domains. 76 /// Domain counters cover all contexts; context filtering happens at each store. 77 struct TeraInvalidationAdmission { 78 private var epoch: String? 79 private var revisions: [TeraRuntimeChangeKind: TeraProjectionRevision] = [:] 80 81 mutating func accept(_ change: TeraRuntimeChange) -> Bool { 82 guard epoch == nil || epoch == change.epoch else { return false } 83 epoch = change.epoch 84 if change.delivery == .resnapshotRequired { 85 return true 86 } 87 if let previous = revisions[change.kind], let next = change.revision.rawValue { 88 guard let value = previous.rawValue, next > value else { return false } 89 } 90 revisions[change.kind] = change.revision 91 return true 92 } 93 } 94 95 struct TeraRuntimeSubscription { 96 let generation: TeraSessionGeneration 97 let continuation: AsyncStream<TeraRuntimeChange>.Continuation 98 var token: (any TeraRuntimeSubscriptionToken)? 99 var admission = TeraInvalidationAdmission() 100 }