Section 15/161 menit
15. Real Use Cases
15. Real Use Cases
Use Case 1: Real-Time Collaboration Editor
Aplikasi collaborative document editing di mana multiple user mengedit dokumen yang sama secara bersamaan, dengan state sync via WebSocket.
swift
// Domain types — semua Sendable untuk safe passing
struct DocumentEdit: Sendable, Codable {
let documentId: UUID
let userId: UUID
let operation: EditOperation
let timestamp: Date
let vectorClock: [UUID: Int]
}
enum EditOperation: Sendable, Codable {
case insert(position: Int, text: String)
case delete(range: Range<Int>)
case replace(range: Range<Int>, text: String)
case format(range: Range<Int>, style: TextStyle)
}
// Core CRDT state — actor-isolated
actor DocumentState {
private var content: AttributedString
private var pendingOperations: [DocumentEdit] = []
private var acknowledgedClock: [UUID: Int] = [:]
init(initialContent: AttributedString) {
self.content = initialContent
}
func applyEdit(_ edit: DocumentEdit) async -> Bool {
// Check vector clock untuk ordering
guard isOrdered(edit.vectorClock) else {
pendingOperations.append(edit)
return false
}
applyOperation(edit.operation)
acknowledgedClock[edit.userId, default: 0] += 1
// Cek apakah ada pending operations yang sekarang bisa diapply
await drainPendingOperations()
return true
}
private func applyOperation(_ operation: EditOperation) {
switch operation {
case .insert(let position, let text):
let idx = content.index(content.startIndex, offsetByCharacters: position)
content.insert(AttributedString(text), at: idx)
case .delete(let range):
let start = content.index(content.startIndex, offsetByCharacters: range.lowerBound)
let end = content.index(content.startIndex, offsetByCharacters: range.upperBound)
content.removeSubrange(start..<end)
case .replace(let range, let text):
let start = content.index(content.startIndex, offsetByCharacters: range.lowerBound)
let end = content.index(content.startIndex, offsetByCharacters: range.upperBound)
content.replaceSubrange(start..<end, with: AttributedString(text))
case .format(let range, let style):
let start = content.index(content.startIndex, offsetByCharacters: range.lowerBound)
let end = content.index(content.startIndex, offsetByCharacters: range.upperBound)
// Apply formatting...
break
}
}
private func drainPendingOperations() async {
var i = 0
while i < pendingOperations.count {
let op = pendingOperations[i]
if isOrdered(op.vectorClock) {
pendingOperations.remove(at: i)
applyOperation(op.operation)
acknowledgedClock[op.userId, default: 0] += 1
} else {
i += 1
}
}
}
private func isOrdered(_ clock: [UUID: Int]) -> Bool {
clock.allSatisfy { userId, count in
acknowledgedClock[userId, default: 0] >= count - 1
}
}
var snapshot: AttributedString { content }
}
// WebSocket connection — actor untuk thread-safe message handling
actor CollaborationSession {
private let documentState: DocumentState
private let currentUserId: UUID
private var outboundQueue: [DocumentEdit] = []
private var webSocketTask: URLSessionWebSocketTask?
init(documentId: UUID, userId: UUID, initialContent: AttributedString) {
self.documentState = DocumentState(initialContent: initialContent)
self.currentUserId = userId
}
func connect(to serverURL: URL) async throws {
let task = URLSession.shared.webSocketTask(with: serverURL)
webSocketTask = task
task.resume()
// Start receive loop
Task { await receiveMessages() }
}
func submitEdit(_ operation: EditOperation) async {
let clock = await currentVectorClock()
let edit = DocumentEdit(
documentId: UUID(),
userId: currentUserId,
operation: operation,
timestamp: Date(),
vectorClock: clock
)
// Apply locally (optimistic update)
await documentState.applyEdit(edit)
// Queue for sending
outboundQueue.append(edit)
await flushOutbound()
}
private func receiveMessages() async {
guard let task = webSocketTask else { return }
while !Task.isCancelled {
do {
let message = try await task.receive()
switch message {
case .data(let data):
let edit = try JSONDecoder().decode(DocumentEdit.self, from: data)
await documentState.applyEdit(edit)
case .string(let text):
guard let data = text.data(using: .utf8),
let edit = try? JSONDecoder().decode(DocumentEdit.self, from: data) else {
continue
}
await documentState.applyEdit(edit)
@unknown default:
break
}
} catch {
// Handle disconnect
break
}
}
}
private func flushOutbound() async {
guard let task = webSocketTask else { return }
for edit in outboundQueue {
guard let data = try? JSONEncoder().encode(edit) else { continue }
try? await task.send(.data(data))
}
outboundQueue.removeAll()
}
private func currentVectorClock() async -> [UUID: Int] {
[currentUserId: 1] // simplified
}
}
Use Case 2: High-Performance Event Processing Pipeline
System event tracking untuk analytics dengan throughput tinggi — harus handle ribuan events per detik tanpa bottleneck.
swift
// Event types — Sendable untuk safe transport
struct AnalyticsEvent: Sendable {
let id: UUID
let name: String
let properties: [String: PropertyValue]
let timestamp: Date
let sessionId: UUID
let userId: UUID?
}
enum PropertyValue: Sendable {
case string(String)
case int(Int)
case double(Double)
case bool(Bool)
}
// High-throughput ingestion dengan batching
actor EventIngestionActor {
private var buffer: [AnalyticsEvent] = []
private let maxBatchSize: Int
private let flushInterval: Duration
private var flushTask: Task<Void, Never>?
private let processingPipeline: EventProcessingPipeline
init(maxBatchSize: Int = 100, flushInterval: Duration = .seconds(5),
pipeline: EventProcessingPipeline) {
self.maxBatchSize = maxBatchSize
self.flushInterval = flushInterval
self.processingPipeline = pipeline
scheduleFlush()
}
func ingest(_ event: AnalyticsEvent) async {
buffer.append(event)
if buffer.count >= maxBatchSize {
await flush()
}
}
private func flush() async {
guard !buffer.isEmpty else { return }
let batch = buffer
buffer.removeAll(keepingCapacity: true)
// Process batch off-actor to avoid blocking ingestion
Task.detached { [batch] in
await self.processingPipeline.process(batch: batch)
}
scheduleFlush()
}
private func scheduleFlush() {
flushTask?.cancel()
flushTask = Task { [weak self, flushInterval] in
try? await Task.sleep(for: flushInterval)
guard !Task.isCancelled else { return }
await self?.flush()
}
}
}
// Pipeline processing — parallelized dengan TaskGroup
actor EventProcessingPipeline {
private let enrichers: [EventEnricher]
private let sinks: [EventSink]
init(enrichers: [EventEnricher], sinks: [EventSink]) {
self.enrichers = enrichers
self.sinks = sinks
}
func process(batch: [AnalyticsEvent]) async {
// Enrich events in parallel
let enriched = await withTaskGroup(of: AnalyticsEvent.self) { group in
for event in batch {
group.addTask { [enrichers = self.enrichers] in
var enrichedEvent = event
for enricher in enrichers {
enrichedEvent = await enricher.enrich(enrichedEvent)
}
return enrichedEvent
}
}
var results: [AnalyticsEvent] = []
for await result in group {
results.append(result)
}
return results
}
// Fan-out to all sinks concurrently
await withTaskGroup(of: Void.self) { group in
for sink in sinks {
let sinkBatch = enriched
group.addTask {
try? await sink.write(events: sinkBatch)
}
}
}
}
}
// Enricher protocol
protocol EventEnricher: Sendable {
func enrich(_ event: AnalyticsEvent) async -> AnalyticsEvent
}
// Sink protocol
protocol EventSink: Sendable {
func write(events: [AnalyticsEvent]) async throws
}
// Concrete enricher — menambah geo-location dari IP
struct GeoEnricher: EventEnricher {
func enrich(_ event: AnalyticsEvent) async -> AnalyticsEvent {
// Lookup geo dari IP — simplified
var enriched = event
// enriched.properties["country"] = .string("ID")
return enriched
}
}
// Metrics untuk monitoring pipeline health — Atomic untuk lock-free
final class PipelineMetrics: Sendable {
let eventsIngested = Atomic<Int>(0)
let eventsProcessed = Atomic<Int>(0)
let processingErrors = Atomic<Int>(0)
let averageLatencyMs = Atomic<Int>(0)
func recordIngestion() {
eventsIngested.fetchAdd(1, ordering: .relaxed)
}
func recordProcessed(count: Int) {
eventsProcessed.fetchAdd(count, ordering: .relaxed)
}
func snapshot() -> (ingested: Int, processed: Int, errors: Int) {
(
eventsIngested.load(ordering: .relaxed),
eventsProcessed.load(ordering: .relaxed),
processingErrors.load(ordering: .relaxed)
)
}
}