Derselbe Fehler wie in @mana/event-sync (dort behoben in 0.15.1), hier in Swift: `sequenceNumber` ist `Int64?`, und SwiftData sortiert über SQLite — dort ist NULL bei aufsteigender Ordnung der KLEINSTE Wert. Frisch emittierte Events (noch ohne Nummer vom Server) landeten damit ganz vorn, und der Replay wandte das Erzeugungs-Event zuletzt an: es überschrieb jede spätere Änderung. Sichtbar in pageta: aus dem Share-Sheet gespeicherte Artikel blieben dauerhaft auf ihrem Platzhalter stehen (Host als Titel, leerer Text), obwohl `enrichPendingArticles()` sauber ArticleEnriched nachgeschickt hatte. In einer echten Datenbank 16 von 22 winfuture-Speicherungen — und winfuture hat gar keine Cookie-Wand, die Extraktion wäre also durchgelaufen. „NULL zuletzt" lässt sich in einem SortDescriptor nicht ausdrücken, deshalb sortiert eventsForAggregate/-Prefix jetzt im Speicher (pro Aggregat wenige Events). Vier Tests decken den gemischten Stream ab; mit dem alten Verhalten fallen drei davon um (verifiziert). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
433 lines
16 KiB
Swift
433 lines
16 KiB
Swift
import Foundation
|
|
import SwiftData
|
|
|
|
// SwiftData-Modelle für event-sync.
|
|
//
|
|
// - `LocalEvent` — Append-only Event-Log. JSON-encoded `EventEnvelope`
|
|
// liegt als `envelopeJSON`; Routing-Felder sind separate Spalten für
|
|
// Index-Abfragen.
|
|
// - `LocalOutboxEntry` — pendende Events, noch nicht an `mana-sync`
|
|
// gepushed. Plus `attempts`/`lastError`.
|
|
// - `LocalMeta` — Key-Value für Cursor + Diagnose + Anon-Identität.
|
|
|
|
@Model
|
|
final class LocalEvent {
|
|
@Attribute(.unique) var idempotencyKey: String
|
|
var eventId: String
|
|
var aggregateId: String
|
|
var appId: String
|
|
var eventType: String
|
|
var attributedToUserId: String
|
|
var sequenceNumber: Int64?
|
|
var occurredAt: String
|
|
var receivedAt: String?
|
|
var envelopeJSON: Data
|
|
|
|
init(
|
|
idempotencyKey: String,
|
|
eventId: String,
|
|
aggregateId: String,
|
|
appId: String,
|
|
eventType: String,
|
|
attributedToUserId: String,
|
|
sequenceNumber: Int64?,
|
|
occurredAt: String,
|
|
receivedAt: String?,
|
|
envelopeJSON: Data
|
|
) {
|
|
self.idempotencyKey = idempotencyKey
|
|
self.eventId = eventId
|
|
self.aggregateId = aggregateId
|
|
self.appId = appId
|
|
self.eventType = eventType
|
|
self.attributedToUserId = attributedToUserId
|
|
self.sequenceNumber = sequenceNumber
|
|
self.occurredAt = occurredAt
|
|
self.receivedAt = receivedAt
|
|
self.envelopeJSON = envelopeJSON
|
|
}
|
|
|
|
func envelope() throws -> EventEnvelope {
|
|
try JSONDecoder().decode(EventEnvelope.self, from: envelopeJSON)
|
|
}
|
|
}
|
|
|
|
@Model
|
|
final class LocalOutboxEntry {
|
|
@Attribute(.unique) var eventId: String
|
|
var idempotencyKey: String
|
|
var envelopeJSON: Data
|
|
var attempts: Int
|
|
var lastError: String?
|
|
var lastAttemptAt: Date?
|
|
var createdAt: Date
|
|
|
|
init(
|
|
eventId: String,
|
|
idempotencyKey: String,
|
|
envelopeJSON: Data,
|
|
attempts: Int = 0,
|
|
lastError: String? = nil,
|
|
lastAttemptAt: Date? = nil,
|
|
createdAt: Date = .init()
|
|
) {
|
|
self.eventId = eventId
|
|
self.idempotencyKey = idempotencyKey
|
|
self.envelopeJSON = envelopeJSON
|
|
self.attempts = attempts
|
|
self.lastError = lastError
|
|
self.lastAttemptAt = lastAttemptAt
|
|
self.createdAt = createdAt
|
|
}
|
|
|
|
func envelope() throws -> EventEnvelope {
|
|
try JSONDecoder().decode(EventEnvelope.self, from: envelopeJSON)
|
|
}
|
|
}
|
|
|
|
@Model
|
|
final class LocalMeta {
|
|
@Attribute(.unique) var key: String
|
|
var value: String
|
|
|
|
init(key: String, value: String) {
|
|
self.key = key
|
|
self.value = value
|
|
}
|
|
}
|
|
|
|
/// Zentraler SwiftData-Container für den Event-Log + Outbox + Meta.
|
|
///
|
|
/// Eine Instanz pro Prozess, `@MainActor`-bound (SwiftData-Standard).
|
|
/// Store-Name, App-Group und app-eigene `@Model`-Typen kommen vom Caller
|
|
/// (die `EventSyncEngine` konstruiert sie aus der `EventSyncConfig`).
|
|
@MainActor
|
|
final class EventStorage {
|
|
let container: ModelContainer
|
|
|
|
init(
|
|
storeName: String,
|
|
appGroupIdentifier: String?,
|
|
additionalModels: [any PersistentModel.Type] = [],
|
|
inMemory: Bool = false
|
|
) throws {
|
|
var models: [any PersistentModel.Type] = [
|
|
LocalEvent.self,
|
|
LocalOutboxEntry.self,
|
|
LocalMeta.self
|
|
]
|
|
models.append(contentsOf: additionalModels)
|
|
let schema = Schema(models)
|
|
|
|
// SwiftData defaultet auf den App-Group-Container, wenn die App
|
|
// ein App-Groups-Entitlement hat. Das Application-Support-Subdir
|
|
// existiert beim ersten Launch noch nicht; wir legen es vorher an,
|
|
// damit SwiftData keinen lauten CoreData-Error-Trace loggt.
|
|
if let appGroupIdentifier,
|
|
let groupURL = FileManager.default.containerURL(
|
|
forSecurityApplicationGroupIdentifier: appGroupIdentifier
|
|
)
|
|
{
|
|
let appSupport = groupURL
|
|
.appendingPathComponent("Library", isDirectory: true)
|
|
.appendingPathComponent("Application Support", isDirectory: true)
|
|
try? FileManager.default.createDirectory(at: appSupport, withIntermediateDirectories: true)
|
|
}
|
|
|
|
let cfg = ModelConfiguration(
|
|
storeName,
|
|
schema: schema,
|
|
isStoredInMemoryOnly: inMemory,
|
|
allowsSave: true
|
|
)
|
|
container = try ModelContainer(for: schema, configurations: [cfg])
|
|
}
|
|
|
|
var context: ModelContext {
|
|
container.mainContext
|
|
}
|
|
|
|
// MARK: - Event-Log
|
|
|
|
@discardableResult
|
|
func appendEvents(_ envelopes: [EventEnvelope]) throws -> Int {
|
|
var inserted = 0
|
|
for envelope in envelopes {
|
|
let key = envelope.idempotencyKey
|
|
var fetch = FetchDescriptor<LocalEvent>(predicate: #Predicate { $0.idempotencyKey == key })
|
|
fetch.fetchLimit = 1
|
|
if try context.fetch(fetch).first != nil { continue }
|
|
|
|
let data = try JSONEncoder().encode(envelope)
|
|
context.insert(LocalEvent(
|
|
idempotencyKey: envelope.idempotencyKey,
|
|
eventId: envelope.eventId,
|
|
aggregateId: envelope.aggregateId,
|
|
appId: envelope.appId,
|
|
eventType: envelope.eventType,
|
|
attributedToUserId: envelope.attributedToUserId,
|
|
sequenceNumber: envelope.sequenceNumber,
|
|
occurredAt: envelope.occurredAt,
|
|
receivedAt: envelope.receivedAt,
|
|
envelopeJSON: data
|
|
))
|
|
inserted += 1
|
|
}
|
|
if inserted > 0 { try context.save() }
|
|
return inserted
|
|
}
|
|
|
|
/// Sortier-Rang eines Events innerhalb seines Aggregats.
|
|
///
|
|
/// Ein Event ohne `sequenceNumber` wurde gerade lokal emittiert und wartet
|
|
/// noch auf seine Nummer vom Server — es ist per Konstruktion **neuer** als
|
|
/// alles Sequenzierte und gehört ans Ende.
|
|
///
|
|
/// Das ist NICHT, was die Datenbank von sich aus tut: SwiftData sortiert
|
|
/// über SQLite, und dort ist `NULL` bei aufsteigender Ordnung der
|
|
/// **kleinste** Wert — lokale Events landeten also ganz vorn. Beim Replay
|
|
/// wurde damit das Erzeugungs-Event (`ArticleSaved` &c.) zuletzt angewandt
|
|
/// und überschrieb jede spätere Änderung.
|
|
///
|
|
/// Sichtbar geworden 2026-07-27 in pageta: aus dem Share-Sheet gespeicherte
|
|
/// Artikel blieben dauerhaft auf ihrem Platzhalter stehen (Host als Titel,
|
|
/// leerer Text), obwohl `enrichPendingArticles()` sauber `ArticleEnriched`
|
|
/// nachgeschickt hatte — 16 von 22 winfuture-Speicherungen in einer
|
|
/// echten Datenbank. Der Web-Client hatte denselben Fehler
|
|
/// (@mana/event-sync ≤ 0.15.0, dort `toBig(null) == 0n`).
|
|
private static func sortRank(_ sequenceNumber: Int64?) -> Int64 {
|
|
sequenceNumber ?? Int64.max
|
|
}
|
|
|
|
/// Events eines Aggregats in Replay-Reihenfolge.
|
|
///
|
|
/// Die Sortierung passiert bewusst **im Speicher**: „NULL zuletzt" lässt
|
|
/// sich in einem `SortDescriptor` nicht ausdrücken, und pro Aggregat sind
|
|
/// es wenige Events.
|
|
func eventsForAggregate(_ aggregateId: String) throws -> [EventEnvelope] {
|
|
let fetch = FetchDescriptor<LocalEvent>(
|
|
predicate: #Predicate { $0.aggregateId == aggregateId }
|
|
)
|
|
return try context.fetch(fetch)
|
|
.sorted { lhs, rhs in
|
|
let l = Self.sortRank(lhs.sequenceNumber)
|
|
let r = Self.sortRank(rhs.sequenceNumber)
|
|
if l != r { return l < r }
|
|
return lhs.occurredAt < rhs.occurredAt
|
|
}
|
|
.map { try $0.envelope() }
|
|
}
|
|
|
|
func eventsForAggregatePrefix(_ prefix: String) throws -> [String: [EventEnvelope]] {
|
|
let needle = prefix + ":"
|
|
let fetch = FetchDescriptor<LocalEvent>(
|
|
predicate: #Predicate { $0.aggregateId.starts(with: needle) }
|
|
)
|
|
var grouped: [String: [LocalEvent]] = [:]
|
|
for row in try context.fetch(fetch) {
|
|
grouped[row.aggregateId, default: []].append(row)
|
|
}
|
|
var out: [String: [EventEnvelope]] = [:]
|
|
for (aggregateId, rows) in grouped {
|
|
out[aggregateId] = try rows
|
|
.sorted { lhs, rhs in
|
|
let l = Self.sortRank(lhs.sequenceNumber)
|
|
let r = Self.sortRank(rhs.sequenceNumber)
|
|
if l != r { return l < r }
|
|
return lhs.occurredAt < rhs.occurredAt
|
|
}
|
|
.map { try $0.envelope() }
|
|
}
|
|
return out
|
|
}
|
|
|
|
func countEvents() throws -> Int {
|
|
try context.fetchCount(FetchDescriptor<LocalEvent>())
|
|
}
|
|
|
|
func clearAll() throws {
|
|
for row in try context.fetch(FetchDescriptor<LocalEvent>()) {
|
|
context.delete(row)
|
|
}
|
|
for row in try context.fetch(FetchDescriptor<LocalOutboxEntry>()) {
|
|
context.delete(row)
|
|
}
|
|
for row in try context.fetch(FetchDescriptor<LocalMeta>()) {
|
|
context.delete(row)
|
|
}
|
|
try context.save()
|
|
}
|
|
|
|
/// Löscht alle lokalen Events, deren `aggregateId` mit `idPrefix` beginnt
|
|
/// (A3-Crypto-Shredding-Begleiter: lokale Klartext-Residue eines gelöschten
|
|
/// Scopes wegräumen). Liefert die Anzahl gelöschter Events.
|
|
@discardableResult
|
|
func purgeByAggregatePrefix(_ idPrefix: String) throws -> Int {
|
|
let rows = try context.fetch(FetchDescriptor<LocalEvent>())
|
|
.filter { $0.aggregateId.hasPrefix(idPrefix) }
|
|
for row in rows { context.delete(row) }
|
|
try context.save()
|
|
return rows.count
|
|
}
|
|
|
|
// MARK: - Claim (Anonymous → echter User)
|
|
|
|
/// Re-attribuiert alle Events + Outbox-Einträge von `oldUserId` auf
|
|
/// `newUserId` (idempotent). Rewritet sowohl die Index-Spalte als
|
|
/// auch das eingebettete Envelope-JSON. Liefert die Anzahl betroffener
|
|
/// Zeilen (Event-Log + Outbox).
|
|
@discardableResult
|
|
func reattribute(from oldUserId: String, to newUserId: String) throws -> Int {
|
|
var touched = 0
|
|
let events = try context.fetch(FetchDescriptor<LocalEvent>(
|
|
predicate: #Predicate { $0.attributedToUserId == oldUserId }
|
|
))
|
|
for row in events {
|
|
let env = try row.envelope().attributed(to: newUserId)
|
|
row.attributedToUserId = newUserId
|
|
row.envelopeJSON = try JSONEncoder().encode(env)
|
|
touched += 1
|
|
}
|
|
let outbox = try context.fetch(FetchDescriptor<LocalOutboxEntry>())
|
|
for row in outbox {
|
|
let env = try row.envelope()
|
|
guard env.attributedToUserId == oldUserId else { continue }
|
|
row.envelopeJSON = try JSONEncoder().encode(env.attributed(to: newUserId))
|
|
touched += 1
|
|
}
|
|
if touched > 0 { try context.save() }
|
|
return touched
|
|
}
|
|
|
|
/// True, wenn mindestens ein Event auf `userId` attribuiert ist. Liest
|
|
/// nur die Index-Spalte (kein JSON-Decode) — billig genug, um bei jedem
|
|
/// Login zu prüfen, ob ein Claim überhaupt nötig ist.
|
|
func hasEvents(attributedTo userId: String) throws -> Bool {
|
|
try context.fetchCount(FetchDescriptor<LocalEvent>(
|
|
predicate: #Predicate { $0.attributedToUserId == userId }
|
|
)) > 0
|
|
}
|
|
|
|
/// Ersetzt die Wire-Envelope eines Outbox-Eintrags (z.B. Re-Encrypt
|
|
/// der Payload beim Claim).
|
|
func updateOutboxEnvelope(eventId: String, _ envelope: EventEnvelope) throws {
|
|
var fetch = FetchDescriptor<LocalOutboxEntry>(predicate: #Predicate { $0.eventId == eventId })
|
|
fetch.fetchLimit = 1
|
|
guard let row = try context.fetch(fetch).first else { return }
|
|
row.envelopeJSON = try JSONEncoder().encode(envelope)
|
|
try context.save()
|
|
}
|
|
|
|
// MARK: - Retention (anonyme Events)
|
|
|
|
/// Begrenzt die anonymen Events (`attributedToUserId` mit `anon:`-Prefix):
|
|
/// löscht zu alte (`maxAgeDays`) und über das Limit hinausgehende
|
|
/// (`maxEvents`, ältestes zuerst) Events plus ihre Outbox-Einträge.
|
|
/// `0` = unbegrenzt. Signierte Events bleiben unangetastet (die liegen
|
|
/// als Cache vor, der Server ist die Wahrheit).
|
|
func pruneAnonymous(maxEvents: Int, maxAgeDays: Int) throws {
|
|
let prefix = "anon:"
|
|
var anon = try context.fetch(FetchDescriptor<LocalEvent>(
|
|
predicate: #Predicate { $0.attributedToUserId.starts(with: prefix) },
|
|
sortBy: [SortDescriptor(\.occurredAt, order: .forward)]
|
|
))
|
|
var toDelete: [LocalEvent] = []
|
|
|
|
if maxAgeDays > 0 {
|
|
let cutoff = ISO8601DateFormatter().string(
|
|
from: Date().addingTimeInterval(-Double(maxAgeDays) * 86400)
|
|
)
|
|
let aged = anon.filter { $0.occurredAt < cutoff }
|
|
toDelete.append(contentsOf: aged)
|
|
let agedIds = Set(aged.map(\.eventId))
|
|
anon.removeAll { agedIds.contains($0.eventId) }
|
|
}
|
|
|
|
if maxEvents > 0, anon.count > maxEvents {
|
|
toDelete.append(contentsOf: anon.prefix(anon.count - maxEvents))
|
|
}
|
|
|
|
guard !toDelete.isEmpty else { return }
|
|
let deletedIds = toDelete.map(\.eventId)
|
|
for row in toDelete {
|
|
context.delete(row)
|
|
}
|
|
try context.save()
|
|
try removeOutbox(eventIds: deletedIds)
|
|
}
|
|
|
|
// MARK: - Outbox
|
|
|
|
func enqueueOutbox(_ envelope: EventEnvelope) throws {
|
|
let eventId = envelope.eventId
|
|
var fetch = FetchDescriptor<LocalOutboxEntry>(predicate: #Predicate { $0.eventId == eventId })
|
|
fetch.fetchLimit = 1
|
|
if try context.fetch(fetch).first != nil { return }
|
|
|
|
try context.insert(LocalOutboxEntry(
|
|
eventId: envelope.eventId,
|
|
idempotencyKey: envelope.idempotencyKey,
|
|
envelopeJSON: JSONEncoder().encode(envelope)
|
|
))
|
|
try context.save()
|
|
}
|
|
|
|
func pendingOutbox(limit: Int = 100) throws -> [LocalOutboxEntry] {
|
|
var fetch = FetchDescriptor<LocalOutboxEntry>(sortBy: [SortDescriptor(\.createdAt, order: .forward)])
|
|
fetch.fetchLimit = limit
|
|
return try context.fetch(fetch)
|
|
}
|
|
|
|
func removeOutbox(eventIds: [String]) throws {
|
|
guard !eventIds.isEmpty else { return }
|
|
let set = Set(eventIds)
|
|
let fetch = FetchDescriptor<LocalOutboxEntry>(predicate: #Predicate { set.contains($0.eventId) })
|
|
for row in try context.fetch(fetch) {
|
|
context.delete(row)
|
|
}
|
|
try context.save()
|
|
}
|
|
|
|
func incrementOutboxAttempt(eventId: String, error: String) throws {
|
|
var fetch = FetchDescriptor<LocalOutboxEntry>(predicate: #Predicate { $0.eventId == eventId })
|
|
fetch.fetchLimit = 1
|
|
guard let row = try context.fetch(fetch).first else { return }
|
|
row.attempts += 1
|
|
row.lastError = error
|
|
row.lastAttemptAt = .init()
|
|
try context.save()
|
|
}
|
|
|
|
func outboxCount() throws -> Int {
|
|
try context.fetchCount(FetchDescriptor<LocalOutboxEntry>())
|
|
}
|
|
|
|
// MARK: - Meta
|
|
|
|
func metaGet(_ key: String) throws -> String? {
|
|
var fetch = FetchDescriptor<LocalMeta>(predicate: #Predicate { $0.key == key })
|
|
fetch.fetchLimit = 1
|
|
return try context.fetch(fetch).first?.value
|
|
}
|
|
|
|
func metaSet(_ key: String, _ value: String) throws {
|
|
var fetch = FetchDescriptor<LocalMeta>(predicate: #Predicate { $0.key == key })
|
|
fetch.fetchLimit = 1
|
|
if let row = try context.fetch(fetch).first {
|
|
row.value = value
|
|
} else {
|
|
context.insert(LocalMeta(key: key, value: value))
|
|
}
|
|
try context.save()
|
|
}
|
|
}
|
|
|
|
enum EventStorageKey {
|
|
static let lastPullCursor = "last-pull-cursor"
|
|
static let lastSyncAt = "last-sync-at"
|
|
static let anonymousUserId = "anonymous-user-id"
|
|
static let clientId = "client-id"
|
|
static let decryptFailureCount = "decrypt-failure-count"
|
|
static let decryptFailureLastAt = "decrypt-failure-last-at"
|
|
static let decryptFailureLastEvent = "decrypt-failure-last-event"
|
|
}
|