EventSyncError.subscriptionRequired (402-Mapping in SyncHTTPClient), onSyncRequired-Hook + isSyncRequired. Push/Pull: Events bleiben in der Outbox (kein Verlust, kein Attempt-Increment), event-getriggerte Push/Pull unterdrückt, Poll prüft weiter (force) → Auto-Recovery. Wire-Parität zu @mana/event-sync 0.9.0. swift build + tests grün. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
136 lines
4.9 KiB
Swift
136 lines
4.9 KiB
Swift
import Foundation
|
|
import ManaCore
|
|
|
|
/// HTTP-Client für die `mana-sync` REST-Surface (`sync2.mana.how`).
|
|
///
|
|
/// Wire siehe `mana/services/mana-sync/internal/handler/sync.go` +
|
|
/// `mana/packages/event-sync/src/client/http.ts`:
|
|
/// - `POST /sync/{appId}` — Append (Batch ≤ 500)
|
|
/// - `GET /sync/{appId}/pull?since=N&limit=K` — Pull ab Cursor
|
|
///
|
|
/// Bearer-Token pro Request frisch aus `AuthClient`. `X-Schema-Hash`
|
|
/// (falls gesetzt) ermöglicht Server-seitige Drift-Erkennung; `422`
|
|
/// signalisiert ein nötiges App-Update (`EventSyncError.schemaOutdated`).
|
|
actor SyncHTTPClient {
|
|
private let baseURL: URL
|
|
private let appId: String
|
|
private let schemaHash: String?
|
|
private let auth: AuthClient
|
|
private let session: URLSession
|
|
private let encoder = JSONEncoder()
|
|
private let decoder = JSONDecoder()
|
|
|
|
init(baseURL: URL, appId: String, schemaHash: String?, auth: AuthClient) {
|
|
self.baseURL = baseURL
|
|
self.appId = appId
|
|
self.schemaHash = schemaHash
|
|
self.auth = auth
|
|
let cfg = URLSessionConfiguration.default
|
|
cfg.timeoutIntervalForRequest = 30
|
|
cfg.timeoutIntervalForResource = 60
|
|
session = URLSession(configuration: cfg)
|
|
}
|
|
|
|
// MARK: - Append
|
|
|
|
struct AppendRequest: Encodable {
|
|
let events: [EventEnvelope]
|
|
}
|
|
|
|
struct AppendedEvent: Decodable {
|
|
let eventId: String
|
|
let sequenceNumber: Int64
|
|
let aggregateId: String
|
|
|
|
// swiftlint:disable:next nesting
|
|
enum CodingKeys: String, CodingKey { case eventId, sequenceNumber, aggregateId }
|
|
|
|
/// sequenceNumber kommt als JSON-String vom Server (`json:",string"`).
|
|
init(from decoder: Decoder) throws {
|
|
let container = try decoder.container(keyedBy: CodingKeys.self)
|
|
eventId = try container.decode(String.self, forKey: .eventId)
|
|
aggregateId = try container.decode(String.self, forKey: .aggregateId)
|
|
sequenceNumber = container.decodeFlexibleInt64(forKey: .sequenceNumber) ?? 0
|
|
}
|
|
}
|
|
|
|
struct RejectedEvent: Decodable {
|
|
let eventId: String
|
|
let errorCode: String
|
|
let humanMessage: String?
|
|
}
|
|
|
|
struct AppendResponse: Decodable {
|
|
let accepted: [AppendedEvent]
|
|
let rejected: [RejectedEvent]
|
|
let nowMs: Int64?
|
|
}
|
|
|
|
func append(_ events: [EventEnvelope]) async throws -> AppendResponse {
|
|
let url = baseURL.appendingPathComponent("sync").appendingPathComponent(appId)
|
|
let body = try encoder.encode(AppendRequest(events: events))
|
|
let (data, http) = try await request(url: url, method: "POST", body: body)
|
|
try ensureOK(http, data: data)
|
|
return try decoder.decode(AppendResponse.self, from: data)
|
|
}
|
|
|
|
// MARK: - Pull
|
|
|
|
struct PullResponse: Decodable {
|
|
let events: [EventEnvelope]
|
|
let hasMore: Bool
|
|
let nextCursor: String?
|
|
}
|
|
|
|
func pull(since: Int64, limit: Int = 100) async throws -> PullResponse {
|
|
var comp = URLComponents(
|
|
url: baseURL
|
|
.appendingPathComponent("sync")
|
|
.appendingPathComponent(appId)
|
|
.appendingPathComponent("pull"),
|
|
resolvingAgainstBaseURL: false
|
|
)!
|
|
comp.queryItems = [
|
|
URLQueryItem(name: "since", value: String(since)),
|
|
URLQueryItem(name: "limit", value: String(limit))
|
|
]
|
|
let (data, http) = try await request(url: comp.url!, method: "GET", body: nil)
|
|
try ensureOK(http, data: data)
|
|
return try decoder.decode(PullResponse.self, from: data)
|
|
}
|
|
|
|
// MARK: - Internals
|
|
|
|
private func request(url: URL, method: String, body: Data?) async throws -> (Data, HTTPURLResponse) {
|
|
var req = URLRequest(url: url)
|
|
req.httpMethod = method
|
|
req.setValue("application/json", forHTTPHeaderField: "Content-Type")
|
|
req.setValue("application/json", forHTTPHeaderField: "Accept")
|
|
if let schemaHash {
|
|
req.setValue(schemaHash, forHTTPHeaderField: "X-Schema-Hash")
|
|
}
|
|
let token = try await auth.freshAccessToken()
|
|
req.setValue("Bearer \(token)", forHTTPHeaderField: "Authorization")
|
|
if let body {
|
|
req.httpBody = body
|
|
}
|
|
let (data, response) = try await session.data(for: req)
|
|
guard let http = response as? HTTPURLResponse else {
|
|
throw EventSyncError.transport(status: -1, body: "invalid response")
|
|
}
|
|
return (data, http)
|
|
}
|
|
|
|
private func ensureOK(_ http: HTTPURLResponse, data: Data) throws {
|
|
if http.statusCode == 422 {
|
|
throw EventSyncError.schemaOutdated
|
|
}
|
|
if http.statusCode == 402 {
|
|
throw EventSyncError.subscriptionRequired
|
|
}
|
|
guard (200 ..< 300).contains(http.statusCode) else {
|
|
let body = String(data: data.prefix(256), encoding: .utf8)
|
|
throw EventSyncError.transport(status: http.statusCode, body: body)
|
|
}
|
|
}
|
|
}
|