androidengineers.Book a session

Data Synchronization

Exercise: Sync Queue with Retries

exercise55 minHard

Build a durable sync queue that persists pending operations in Room, processes them in order via WorkManager, and retries failed operations with exponential backoff.

The Outbox Table

@Entity(tableName = "sync_queue")
data class SyncOperation(
    @PrimaryKey(autoGenerate = true) val id: Long = 0,
    val type: String,          // "create_post", "update_post", "delete_post"
    val entityId: String,
    val payload: String,       // JSON body to send to API
    val priority: Int = 0,     // higher = processed first
    val retryCount: Int = 0,
    val lastAttemptAt: Long? = null,
    val createdAt: Long = System.currentTimeMillis(),
    val status: String = "PENDING"  // PENDING, PROCESSING, FAILED
)

@Dao
interface SyncQueueDao {
    @Insert suspend fun enqueue(op: SyncOperation): Long
    @Query("SELECT * FROM sync_queue WHERE status = 'PENDING' ORDER BY priority DESC, createdAt ASC LIMIT 10")
    suspend fun nextBatch(): List<SyncOperation>
    @Update suspend fun update(op: SyncOperation)
    @Delete suspend fun remove(op: SyncOperation)
    @Query("SELECT COUNT(*) FROM sync_queue WHERE status = 'PENDING'") suspend fun pendingCount(): Int
}

Step 1: Enqueueing Operations

class SyncQueueManager(
    private val dao: SyncQueueDao,
    private val context: Context
) {
    suspend fun enqueue(type: String, entityId: String, payload: String) {
        dao.enqueue(SyncOperation(type = type, entityId = entityId, payload = payload))
        scheduleProcessing()
    }

    private fun scheduleProcessing() {
        val request = OneTimeWorkRequestBuilder<SyncQueueWorker>()
            .setConstraints(
                Constraints.Builder().setRequiredNetworkType(NetworkType.CONNECTED).build()
            )
            .setBackoffCriteria(BackoffPolicy.EXPONENTIAL, 30, TimeUnit.SECONDS)
            .build()

        WorkManager.getInstance(context).enqueueUniqueWork(
            "sync_queue_processing",
            ExistingWorkPolicy.KEEP,  // don't interrupt in-progress processing
            request
        )
    }
}

Step 2: Processing Worker

class SyncQueueWorker(context: Context, params: WorkerParameters) : CoroutineWorker(context, params) {
    private val dao = AppDatabase.getInstance(context).syncQueueDao()
    private val api = RetrofitClient.api

    override suspend fun doWork(): Result {
        val operations = dao.nextBatch()
        if (operations.isEmpty()) return Result.success()

        var hasFailures = false

        for (op in operations) {
            // Mark as processing
            dao.update(op.copy(status = "PROCESSING", lastAttemptAt = System.currentTimeMillis()))

            val success = try {
                processOperation(op)
                dao.remove(op)
                true
            } catch (e: IOException) {
                hasFailures = true
                val maxRetries = 5
                if (op.retryCount >= maxRetries) {
                    dao.update(op.copy(status = "FAILED"))
                } else {
                    dao.update(op.copy(status = "PENDING", retryCount = op.retryCount + 1))
                }
                false
            } catch (e: HttpException) {
                if (e.code() in 400..499) {
                    // Client error — don't retry, mark as failed
                    dao.update(op.copy(status = "FAILED"))
                } else {
                    hasFailures = true
                    dao.update(op.copy(status = "PENDING", retryCount = op.retryCount + 1))
                }
                false
            }
        }

        // If there are more pending items, schedule another run
        return if (dao.pendingCount() > 0 || hasFailures) Result.retry() else Result.success()
    }

    private suspend fun processOperation(op: SyncOperation) {
        when (op.type) {
            "create_post" -> {
                val post = json.decodeFromString<Post>(op.payload)
                api.createPost(post)
            }
            "update_post" -> {
                val post = json.decodeFromString<Post>(op.payload)
                api.updatePost(post.id, post)
            }
            "delete_post" -> api.deletePost(op.entityId)
            else -> throw IllegalArgumentException("Unknown operation type: ${op.type}")
        }
    }
}

Step 3: Triggering on Connectivity Change

class ConnectivityObserver(context: Context) {
    val isConnected: Flow<Boolean> = callbackFlow {
        val cm = context.getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager
        val callback = object : ConnectivityManager.NetworkCallback() {
            override fun onAvailable(network: Network) { trySend(true) }
            override fun onLost(network: Network) { trySend(false) }
        }
        cm.registerDefaultNetworkCallback(callback)
        awaitClose { cm.unregisterNetworkCallback(callback) }
    }
}

// In Application.onCreate:
lifecycleScope.launch {
    ConnectivityObserver(context).isConnected
        .filter { it }  // only when connected
        .collect {
            syncQueueManager.scheduleProcessing()
        }
}

Step 4: Unit Test

class SyncQueueWorkerTest {
    private val fakeDao = FakeSyncQueueDao()
    private val fakeApi = FakeApi()

    @Test fun `successful operation is removed from queue`() = runTest {
        fakeDao.enqueue(SyncOperation(type = "create_post", entityId = "p1", payload = "{}"))

        val worker = SyncQueueWorker(/* ... */)
        val result = worker.doWork()

        assertEquals(Result.success(), result)
        assertEquals(0, fakeDao.pendingCount())
    }

    @Test fun `retryable failure increments retryCount`() = runTest {
        fakeDao.enqueue(SyncOperation(type = "create_post", entityId = "p1", payload = "{}"))
        fakeApi.shouldThrow = IOException("Network error")

        val worker = SyncQueueWorker(/* ... */)
        worker.doWork()

        val op = fakeDao.nextBatch().first()
        assertEquals(1, op.retryCount)
        assertEquals("PENDING", op.status)
    }
}

Key Takeaways

ConceptRule
ExistingWorkPolicy.KEEPPrevents duplicate processing workers
4xx errorsDon't retry client errors — mark as FAILED
5xx / IOExceptionRetry with exponential backoff
pendingCount() > 0Return retry() to process remaining items
Connectivity triggerSchedule processing when network becomes available

YOUR LEARNING JOURNEY

0 of 177 available lessons completed

Progress saved in this browser. No account needed.
Exercise: Sync Queue with Retries | Android System Design | Android Engineers