Section 14/211 menit

14. Swift Concurrency: AsyncStream — Live Polling

14. Swift Concurrency: AsyncStream — Live Polling

Problem

Beberapa layar membutuhkan data yang terus-menerus diperbarui — polling setiap N detik, atau menerima push dari server. Dengan callback/timer, kode menjadi tersebar dan sulit di-cancel dengan bersih. AsyncStream memberikan model push-based yang bisa di-for await.

Konsep

AsyncStream<Element> adalah sequence async yang menghasilkan nilai satu per satu via continuation.yield(). Producer (yang menghasilkan data) dan consumer (yang mengkonsumsi) berjalan secara independen — producer mengirim kapan pun siap, consumer menerima saat ada nilai baru.

Implementasi Live Polling di Worker

swift
actor UserWorker {
    // ...existing methods...

    // Menghasilkan [User] secara berkala — setiap 'interval' sekali.
    // Caller cukup `for await users in worker.liveUserStream()` — polling
    // ditangani sepenuhnya di sini, termasuk cleanup saat Task di-cancel.
    func liveUserStream(interval: Duration = .seconds(30)) -> AsyncStream<[User]> {
        AsyncStream { continuation in
            let task = Task {
                // Fetch pertama segera — jangan tunggu interval pertama.
                if let users = try? await fetchUsers() {
                    continuation.yield(users)
                }

                // Loop sampai Task di-cancel.
                while !Task.isCancelled {
                    // Tidur dulu sebelum fetch berikutnya.
                    // Task.sleep throws CancellationError jika Task di-cancel saat sleeping.
                    try? await Task.sleep(for: interval)

                    // Cek lagi setelah sleep — siapa tahu di-cancel persis saat tidur.
                    guard !Task.isCancelled else { break }

                    if let users = try? await fetchUsers() {
                        continuation.yield(users)
                    }
                }

                // Sinyal bahwa stream sudah selesai.
                continuation.finish()
            }

            // Ketika consumer tidak lagi mengkonsumsi stream (mis. view dismissed),
            // onTermination dipanggil otomatis — kita cancel polling task.
            continuation.onTermination = { _ in
                task.cancel()
            }
        }
    }
}

Implementasi di Interactor

swift
actor UserListInteractor: UserListBusinessLogic, UserListDataStore {
    // ...existing properties...

    // Simpan task polling agar bisa di-cancel dari luar.
    private var liveUpdateTask: Task<Void, Never>?

    func startLiveUpdates() {
        // Cegah multiple polling task berjalan bersamaan.
        liveUpdateTask?.cancel()

        liveUpdateTask = Task {
            // AsyncStream dari Worker — setiap 30 detik ada nilai baru.
            // 'for await' di sini adalah loop yang terus berjalan
            // sampai stream selesai atau Task di-cancel.
            for await users in worker.liveUserStream(interval: .seconds(30)) {
                // Guard penting: cek cancellation di setiap iterasi.
                // Task.isCancelled bisa true di antara dua iterasi.
                guard !Task.isCancelled else { break }

                cachedUsers = users
                let response = UserList.FetchUsers.Response(users: users, stats: nil, featuredUser: nil)
                await presenter?.presentUsers(response)
            }
        }
    }

    func stopLiveUpdates() {
        // Cancel task polling → onTermination dipanggil → inner Task di Worker di-cancel
        // → loop di Worker berhenti → stream selesai → for await di Interactor berhenti.
        // Ini adalah chain cancellation yang bersih dan terstruktur.
        liveUpdateTask?.cancel()
        liveUpdateTask = nil
    }
}

Implementasi di ViewController

swift
@MainActor
final class UserListViewController: UIViewController {
    var interactor: (any UserListBusinessLogic)?
    var router: (any UserListRoutingLogic)?

    override func viewDidAppear(_ animated: Bool) {
        super.viewDidAppear(animated)
        // Mulai live polling saat layar terlihat.
        Task { await interactor?.startLiveUpdates() }
    }

    override func viewDidDisappear(_ animated: Bool) {
        super.viewDidDisappear(animated)
        // Hentikan polling saat layar tidak terlihat — hemat baterai dan bandwidth.
        Task { await interactor?.stopLiveUpdates() }
    }
}

Tambahkan ke Protocol

swift
protocol UserListBusinessLogic: AnyObject, Sendable {
    func fetchUsers(request: UserList.FetchUsers.Request) async
    func selectUser(request: UserList.SelectUser.Request) async
    func startLiveUpdates() async    // tambahan
    func stopLiveUpdates() async     // tambahan
}

Variasi: AsyncStream dari Notification

Scenario: refresh list ketika app kembali dari background.

swift
extension NotificationCenter {
    // Wrap NotificationCenter menjadi AsyncStream — tidak perlu addObserver/removeObserver manual.
    func notifications(named name: Notification.Name) -> AsyncStream<Notification> {
        AsyncStream { continuation in
            let observer = addObserver(forName: name, object: nil, queue: nil) { notification in
                continuation.yield(notification)
            }
            continuation.onTermination = { [weak self] _ in
                self?.removeObserver(observer)
            }
        }
    }
}

// Penggunaan di Interactor:
func observeAppForeground() {
    Task {
        for await _ in NotificationCenter.default.notifications(
            named: UIApplication.willEnterForegroundNotification
        ) {
            // Refresh saat app kembali foreground
            guard !Task.isCancelled else { break }
            await fetchUsers(request: UserList.FetchUsers.Request())
        }
    }
}

Kapan Pakai AsyncStream:

  • Polling berkala yang harus bisa di-cancel dengan bersih
  • Mengubah event-based API (timer, NotificationCenter, delegate) menjadi async sequence
  • Real-time updates yang datang secara push (WebSocket, server-sent events)

Kapan Tidak Pakai AsyncStream:

  • Operasi satu kali (sekali fetch, satu response) → gunakan async/await biasa
  • Sudah punya Combine publisher yang reliable di codebase → pertimbangkan konsistensi