mana-swift-event-sync/Sources/ManaEventSync/Storage/EventStorage.swift
Till JS b81ffb4c0b fix: lokal emittierte Events sortieren nicht mehr vor den gesyncten
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>
2026-07-27 14:52:37 +02:00

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"
}