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