Section 5/201 menit

5. AsyncSequence & AsyncStream

5. AsyncSequence & AsyncStream

Teori: Push vs Pull dan Backpressure

Ada dua model fundamental untuk streaming data:

Push-based (Combine, NotificationCenter, delegates):

swift
Producer ──push──▶ Consumer
"Data tersedia? Saya kirim sekarang"

Masalah: Producer bisa lebih cepat dari Consumer → Consumer kewalahan (perlu backpressure)

Pull-based (AsyncSequence):

swift
Consumer ──request──▶ Producer ──yield──▶ Consumer
"Saya siap menerima data berikutnya"

AsyncSequence menggunakan model pull: consumer meminta elemen berikutnya, dan hanya menerimanya ketika siap. Ini menciptakan backpressure natural — producer tidak bisa mengirim lebih cepat dari consumer bisa proses.

AsyncSequence — Protocol

swift
// Definisi minimal AsyncSequence:
protocol AsyncSequence {
    associatedtype Element
    associatedtype AsyncIterator: AsyncIteratorProtocol

    func makeAsyncIterator() -> AsyncIterator
}

protocol AsyncIteratorProtocol {
    associatedtype Element
    mutating func next() async throws -> Element?
}

Penggunaan dengan for await:

swift
// for await adalah sintaks khusus untuk AsyncSequence
for await notification in NotificationCenter.default.notifications(named: .UIKeyboardWillShow) {
    let frame = notification.userInfo?[UIKeyboardFrameEndUserInfoKey] as? CGRect
    adjustForKeyboard(frame: frame)
}

// Dengan filter, map, dll — AsyncSequence punya lazy operators
let recentErrors = logStream
    .filter { $0.level == .error }
    .prefix(10)  // hanya 10 pertama

AsyncStream — Bridge dari Callback ke Async

AsyncStream adalah cara membuat AsyncSequence dari sumber data push-based:

swift
// Contoh: wrapping NotificationCenter ke AsyncStream
func keyboardNotifications() -> AsyncStream<KeyboardInfo> {
    AsyncStream { continuation in
        let showObserver = NotificationCenter.default.addObserver(
            forName: UIResponder.keyboardWillShowNotification,
            object: nil, queue: nil
        ) { notification in
            let info = KeyboardInfo(from: notification, isShowing: true)
            continuation.yield(info)  // push elemen ke stream
        }

        let hideObserver = NotificationCenter.default.addObserver(
            forName: UIResponder.keyboardWillHideNotification,
            object: nil, queue: nil
        ) { notification in
            let info = KeyboardInfo(from: notification, isShowing: false)
            continuation.yield(info)
        }

        // Cleanup saat stream selesai/di-cancel
        continuation.onTermination = { _ in
            NotificationCenter.default.removeObserver(showObserver)
            NotificationCenter.default.removeObserver(hideObserver)
        }
    }
}

// Penggunaan:
Task { @MainActor in
    for await keyboard in keyboardNotifications() {
        adjustLayout(for: keyboard)
    }
}

AsyncThrowingStream — Dengan Error Handling

swift
func liveUpdates(for documentID: String) -> AsyncThrowingStream<Document, Error> {
    AsyncThrowingStream { continuation in
        let listener = FirebaseDB.collection("docs").document(documentID)
            .addSnapshotListener { snapshot, error in
                if let error = error {
                    continuation.finish(throwing: error)  // stream selesai dengan error
                    return
                }
                guard let doc = try? snapshot?.data(as: Document.self) else {
                    return  // skip invalid snapshot
                }
                continuation.yield(doc)  // push document update
            }

        continuation.onTermination = { _ in
            listener.remove()  // cleanup Firebase listener
        }
    }
}

// Penggunaan:
do {
    for try await document in liveUpdates(for: "doc_123") {
        updateUI(with: document)
    }
} catch {
    showError(error)
}

Kapan Menggunakan

AsyncSequence — gunakan untuk:

  • Streaming data: WebSocket messages, server-sent events
  • Real-time updates: database listeners, file watchers
  • Infinite sequences: timer ticks, sensor data
  • Mengganti delegation pattern yang berulang

AsyncStream — gunakan ketika:

  • Menjembatani callback/delegate API ke async context
  • Wrapping notification-based API
  • Membuat custom event emitter

Jangan gunakan ketika:

  • Hanya butuh satu nilai async (gunakan async/await biasa)
  • Kamu butuh backpressure yang sangat presisi (gunakan custom AsyncIterator)
  • Data sudah tersedia sepenuhnya (gunakan Array biasa)

Real Use Case: WebSocket Live Chat

swift
actor WebSocketClient {
    enum Message: Sendable {
        case text(String)
        case data(Data)
        case connected
        case disconnected(Error?)
    }

    private var webSocketTask: URLSessionWebSocketTask?
    private var continuation: AsyncStream<Message>.Continuation?

    var messages: AsyncStream<Message> {
        AsyncStream { [weak self] continuation in
            self?.continuation = continuation
            continuation.onTermination = { [weak self] _ in
                Task { await self?.disconnect() }
            }
        }
    }

    func connect(to url: URL) async throws {
        let task = URLSession.shared.webSocketTask(with: url)
        webSocketTask = task
        task.resume()
        continuation?.yield(.connected)
        await receiveLoop()
    }

    private func receiveLoop() async {
        while let task = webSocketTask {
            do {
                let message = try await task.receive()
                switch message {
                case .string(let text): continuation?.yield(.text(text))
                case .data(let data): continuation?.yield(.data(data))
                @unknown default: break
                }
            } catch {
                continuation?.yield(.disconnected(error))
                continuation?.finish()
                break
            }
        }
    }

    func send(_ text: String) async throws {
        try await webSocketTask?.send(.string(text))
    }

    func disconnect() {
        webSocketTask?.cancel(with: .normalClosure, reason: nil)
        webSocketTask = nil
        continuation?.finish()
        continuation = nil
    }
}

// Penggunaan di ViewModel:
@MainActor
class ChatViewModel: ObservableObject {
    @Published var messages: [ChatMessage] = []
    private let client = WebSocketClient()
    private var receiveTask: Task<Void, Never>?

    func connect(to url: URL) async throws {
        try await client.connect(to: url)

        receiveTask = Task {
            for await message in await client.messages {
                switch message {
                case .text(let text):
                    messages.append(ChatMessage(content: text, isMine: false))
                case .disconnected:
                    showReconnectPrompt()
                default: break
                }
            }
        }
    }

    func send(text: String) async throws {
        messages.append(ChatMessage(content: text, isMine: true))
        try await client.send(text)
    }
}

struct ChatMessage: Identifiable {
    let id = UUID()
    let content: String
    let isMine: Bool
}

SWIFT 6.0 — STRICT CONCURRENCY (2024)