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
| Concept | Rule |
|---|---|
ExistingWorkPolicy.KEEP | Prevents duplicate processing workers |
| 4xx errors | Don't retry client errors — mark as FAILED |
| 5xx / IOException | Retry with exponential backoff |
pendingCount() > 0 | Return retry() to process remaining items |
| Connectivity trigger | Schedule processing when network becomes available |