Implement a reliable message send pipeline: optimistic insert, background retry on failure, WebSocket ack handling, and UI status indicators for each message state.
Goal
A ConversationScreen where:
- Messages appear instantly when sent (optimistic)
- Failed messages show a retry option
- Sent messages show checkmarks (sent → delivered → read)
- Background retry using WorkManager survives app restart
Step 1: Database Schema
@Entity(tableName = "messages")
data class MessageEntity(
@PrimaryKey val id: String,
val conversationId: String,
val body: String,
val senderId: String,
val status: MessageStatus,
val createdAt: Long,
val localCreatedAt: Long,
val localId: String, // client-generated; idempotency key
val retryCount: Int = 0,
val isOutgoing: Boolean
)
enum class MessageStatus { SENDING, SENT, DELIVERED, READ, FAILED }
@Dao
interface MessageDao {
@Query("SELECT * FROM messages WHERE conversation_id = :id ORDER BY local_created_at ASC")
fun observeMessages(id: String): Flow<List<MessageEntity>>
@Insert(onConflict = OnConflictStrategy.REPLACE)
suspend fun upsert(message: MessageEntity)
@Query("UPDATE messages SET status = :status WHERE local_id = :localId")
suspend fun updateStatus(localId: String, status: MessageStatus)
@Query("UPDATE messages SET retry_count = retry_count + 1 WHERE local_id = :localId")
suspend fun incrementRetry(localId: String)
@Query("SELECT * FROM messages WHERE status = 'FAILED' AND is_outgoing = 1")
suspend fun getFailedMessages(): List<MessageEntity>
}
Step 2: Send Pipeline
class MessageSendUseCase @Inject constructor(
private val dao: MessageDao,
private val api: MessagingApi,
private val workManager: WorkManager
) {
suspend fun send(conversationId: String, body: String) {
val localId = UUID.randomUUID().toString()
val now = System.currentTimeMillis()
// Step 1: Optimistic insert
val message = MessageEntity(
id = localId, // temp ID = local ID initially
conversationId = conversationId,
body = body,
senderId = currentUserId,
status = MessageStatus.SENDING,
createdAt = now,
localCreatedAt = now,
localId = localId,
isOutgoing = true
)
dao.upsert(message)
// Step 2: Attempt send
try {
val response = api.sendMessage(
conversationId = conversationId,
body = body,
idempotencyKey = localId
)
// Step 3a: Success — update with server ID and SENT status
dao.upsert(message.copy(
id = response.messageId,
status = MessageStatus.SENT,
createdAt = response.serverTimestamp
))
} catch (e: IOException) {
// Step 3b: Network failure — mark FAILED, schedule retry
dao.updateStatus(localId, MessageStatus.FAILED)
scheduleRetry(localId, conversationId, body)
}
}
private fun scheduleRetry(localId: String, conversationId: String, body: String) {
val work = OneTimeWorkRequestBuilder<MessageRetryWorker>()
.setInputData(workDataOf(
"local_id" to localId,
"conversation_id" to conversationId,
"body" to body
))
.setConstraints(Constraints.Builder()
.setRequiredNetworkType(NetworkType.CONNECTED)
.build())
.setBackoffCriteria(BackoffPolicy.EXPONENTIAL, 15, TimeUnit.SECONDS)
.addTag("message_retry")
.build()
workManager.enqueueUniqueWork(
"retry_$localId",
ExistingWorkPolicy.KEEP, // don't re-enqueue if already pending
work
)
}
}
Step 3: Retry Worker
class MessageRetryWorker(
context: Context,
params: WorkerParameters,
private val api: MessagingApi,
private val dao: MessageDao
) : CoroutineWorker(context, params) {
override suspend fun doWork(): Result {
val localId = inputData.getString("local_id") ?: return Result.failure()
val conversationId = inputData.getString("conversation_id") ?: return Result.failure()
val body = inputData.getString("body") ?: return Result.failure()
dao.incrementRetry(localId)
dao.updateStatus(localId, MessageStatus.SENDING)
return try {
val response = api.sendMessage(conversationId, body, idempotencyKey = localId)
val currentMessage = dao.getByLocalId(localId) ?: return Result.failure()
dao.upsert(currentMessage.copy(
id = response.messageId,
status = MessageStatus.SENT,
createdAt = response.serverTimestamp
))
Result.success()
} catch (e: IOException) {
dao.updateStatus(localId, MessageStatus.FAILED)
if (runAttemptCount < 5) Result.retry() else Result.failure()
}
}
}
Step 4: WebSocket Acks
class MessageAckHandler(
private val socket: MessagingWebSocket,
private val dao: MessageDao
) {
fun startHandling(): Job = coroutineScope.launch {
socket.events.collect { event ->
when (event) {
is SocketEvent.MessageDelivered ->
dao.updateStatus(event.messageLocalId, MessageStatus.DELIVERED)
is SocketEvent.MessageRead ->
dao.updateStatus(event.messageLocalId, MessageStatus.READ)
is SocketEvent.IncomingMessage ->
dao.upsert(event.message.toEntity())
}
}
}
}
Step 5: UI
@Composable
fun MessageBubble(message: Message, onRetry: (String) -> Unit) {
Row(
horizontalArrangement = if (message.isOutgoing) Arrangement.End else Arrangement.Start,
modifier = Modifier.fillMaxWidth().padding(horizontal = 8.dp, vertical = 2.dp)
) {
Column(horizontalAlignment = if (message.isOutgoing) Alignment.End else Alignment.Start) {
Surface(
shape = RoundedCornerShape(12.dp),
color = if (message.isOutgoing) MaterialTheme.colorScheme.primary
else MaterialTheme.colorScheme.surfaceVariant
) {
Text(message.body, modifier = Modifier.padding(8.dp))
}
// Status indicator
Row(verticalAlignment = Alignment.CenterVertically) {
when (message.status) {
MessageStatus.SENDING ->
CircularProgressIndicator(modifier = Modifier.size(12.dp))
MessageStatus.SENT ->
Icon(Icons.Default.Check, null, modifier = Modifier.size(12.dp))
MessageStatus.DELIVERED ->
Icon(Icons.Default.DoneAll, null, modifier = Modifier.size(12.dp))
MessageStatus.READ ->
Icon(Icons.Default.DoneAll, null,
modifier = Modifier.size(12.dp),
tint = MaterialTheme.colorScheme.primary)
MessageStatus.FAILED ->
TextButton(onClick = { onRetry(message.localId) }) {
Text("Retry", color = MaterialTheme.colorScheme.error, fontSize = 11.sp)
}
}
}
}
}
}
Verification Checklist
[ ] Send message → appears immediately with spinner
[ ] Spinner disappears on SENT status (server ack)
[ ] Check mark updates when recipient delivers/reads
[ ] Force kill + reopen: FAILED messages still shown with Retry
[ ] Tap Retry → message transitions FAILED → SENDING → SENT
[ ] Network off → all sends go FAILED; WorkManager retries when network returns
[ ] No duplicate messages after retry (idempotency key prevents double-send)