diff --git a/apps/.i18n/native-source.json b/apps/.i18n/native-source.json index 7b52418a975c..93c38cca7133 100644 --- a/apps/.i18n/native-source.json +++ b/apps/.i18n/native-source.json @@ -5819,7 +5819,7 @@ }, { "kind": "ui-named-argument", - "line": 436, + "line": 437, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt", "source": "You", "surface": "android", @@ -5827,7 +5827,15 @@ }, { "kind": "ui-named-argument", - "line": 450, + "line": 444, + "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt", + "source": "📎 ${attachment.fileName}", + "surface": "android", + "id": "native.android.10d8391dea9f334b" + }, + { + "kind": "ui-named-argument", + "line": 460, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt", "source": "Retry", "surface": "android", @@ -5835,7 +5843,7 @@ }, { "kind": "ui-named-argument", - "line": 453, + "line": 465, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt", "source": "Delete", "surface": "android", @@ -5843,7 +5851,7 @@ }, { "kind": "ui-named-argument", - "line": 485, + "line": 497, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt", "source": "OpenClaw · Live", "surface": "android", @@ -5851,7 +5859,7 @@ }, { "kind": "conditional-branch", - "line": 522, + "line": 534, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt", "source": "System", "surface": "android", @@ -5859,7 +5867,7 @@ }, { "kind": "conditional-branch", - "line": 523, + "line": 535, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt", "source": "OpenClaw", "surface": "android", @@ -5867,7 +5875,7 @@ }, { "kind": "ui-call", - "line": 549, + "line": 561, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt", "source": "Unsupported attachment", "surface": "android", @@ -5883,7 +5891,7 @@ }, { "kind": "ui-named-argument", - "line": 517, + "line": 518, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "All", "surface": "android", @@ -5891,7 +5899,7 @@ }, { "kind": "ui-named-argument", - "line": 575, + "line": 576, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "OpenClaw", "surface": "android", @@ -5899,7 +5907,7 @@ }, { "kind": "ui-named-argument", - "line": 596, + "line": 597, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "New chat", "surface": "android", @@ -5907,7 +5915,7 @@ }, { "kind": "ui-named-argument", - "line": 601, + "line": 602, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "More new chat options", "surface": "android", @@ -5915,7 +5923,7 @@ }, { "kind": "ui-named-argument", - "line": 616, + "line": 617, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Refresh chat", "surface": "android", @@ -5923,7 +5931,7 @@ }, { "kind": "ui-named-argument", - "line": 619, + "line": 620, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Chat", "surface": "android", @@ -5931,7 +5939,7 @@ }, { "kind": "ui-named-argument", - "line": 767, + "line": 768, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Loading session", "surface": "android", @@ -5939,7 +5947,7 @@ }, { "kind": "ui-named-argument", - "line": 794, + "line": 795, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Jump to latest", "surface": "android", @@ -5947,7 +5955,7 @@ }, { "kind": "conditional-branch", - "line": 820, + "line": 821, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Ready when you are", "surface": "android", @@ -5955,7 +5963,7 @@ }, { "kind": "ui-named-argument", - "line": 851, + "line": 852, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Fix connection", "surface": "android", @@ -5963,7 +5971,7 @@ }, { "kind": "ui-named-argument", - "line": 852, + "line": 853, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Copy diagnostics", "surface": "android", @@ -5971,7 +5979,7 @@ }, { "kind": "conditional-branch", - "line": 1036, + "line": 1037, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Preparing audio…", "surface": "android", @@ -5979,7 +5987,7 @@ }, { "kind": "conditional-branch", - "line": 1036, + "line": 1037, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Speaking…", "surface": "android", @@ -5987,7 +5995,7 @@ }, { "kind": "ui-named-argument", - "line": 1057, + "line": 1058, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Tools running", "surface": "android", @@ -5995,7 +6003,7 @@ }, { "kind": "ui-named-argument", - "line": 1059, + "line": 1060, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "OpenClaw is working", "surface": "android", @@ -6003,7 +6011,7 @@ }, { "kind": "ui-named-argument", - "line": 1062, + "line": 1063, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "+${toolCalls.size - 4} more", "surface": "android", @@ -6011,7 +6019,7 @@ }, { "kind": "ui-named-argument", - "line": 1072, + "line": 1073, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Thinking", "surface": "android", @@ -6019,7 +6027,7 @@ }, { "kind": "ui-named-argument", - "line": 1073, + "line": 1074, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "OpenClaw is preparing a response.", "surface": "android", @@ -6027,7 +6035,7 @@ }, { "kind": "ui-named-argument", - "line": 1255, + "line": 1253, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Stop", "surface": "android", @@ -6035,7 +6043,7 @@ }, { "kind": "ui-named-argument", - "line": 1352, + "line": 1350, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Default", "surface": "android", @@ -6043,7 +6051,7 @@ }, { "kind": "conditional-branch", - "line": 1418, + "line": 1416, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Pin model", "surface": "android", @@ -6051,7 +6059,7 @@ }, { "kind": "conditional-branch", - "line": 1418, + "line": 1416, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Unpin model", "surface": "android", @@ -6059,7 +6067,7 @@ }, { "kind": "ui-named-argument", - "line": 1435, + "line": 1433, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "No commands found", "surface": "android", @@ -6067,7 +6075,7 @@ }, { "kind": "ui-named-argument", - "line": 1497, + "line": 1495, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Gateway offline", "surface": "android", @@ -6075,7 +6083,7 @@ }, { "kind": "conditional-branch", - "line": 1543, + "line": 1541, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Close thinking level selector", "surface": "android", @@ -6083,7 +6091,7 @@ }, { "kind": "conditional-branch", - "line": 1543, + "line": 1541, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Open thinking level selector", "surface": "android", @@ -6091,7 +6099,7 @@ }, { "kind": "ui-named-argument", - "line": 1603, + "line": 1601, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Attach image", "surface": "android", @@ -6099,7 +6107,7 @@ }, { "kind": "ui-named-argument", - "line": 1638, + "line": 1636, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Message OpenClaw", "surface": "android", @@ -6107,7 +6115,7 @@ }, { "kind": "ui-named-argument", - "line": 1654, + "line": 1652, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Open voice", "surface": "android", @@ -6115,7 +6123,7 @@ }, { "kind": "ui-named-argument", - "line": 1703, + "line": 1701, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Remove attachment", "surface": "android", @@ -6123,7 +6131,7 @@ }, { "kind": "conditional-branch", - "line": 1724, + "line": 1722, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Main", "surface": "android", @@ -6131,7 +6139,7 @@ }, { "kind": "conditional-branch", - "line": 1732, + "line": 1730, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "$emoji $name", "surface": "android", @@ -6139,7 +6147,7 @@ }, { "kind": "ui-named-argument", - "line": 1791, + "line": 1789, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "Send", "surface": "android", @@ -6147,7 +6155,7 @@ }, { "kind": "conditional-branch", - "line": 1822, + "line": 1820, "path": "apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt", "source": "$contextLabel · ${contextMeterThinkingLabel(thinkingLevel)}", "surface": "android", diff --git a/apps/android/app/src/main/java/ai/openclaw/app/NodeApp.kt b/apps/android/app/src/main/java/ai/openclaw/app/NodeApp.kt index 36d6fc687166..246277bd3341 100644 --- a/apps/android/app/src/main/java/ai/openclaw/app/NodeApp.kt +++ b/apps/android/app/src/main/java/ai/openclaw/app/NodeApp.kt @@ -1,6 +1,7 @@ package ai.openclaw.app import ai.openclaw.app.chat.ChatCacheDatabase +import ai.openclaw.app.chat.RoomChatCommandOutbox import ai.openclaw.app.gateway.DeviceAuthStore import ai.openclaw.app.gateway.DeviceIdentityStore import android.app.Application @@ -60,7 +61,8 @@ class NodeApp : Application() { database.withTransaction { database.dao().deleteMessages(gatewayId) database.dao().deleteSessions(gatewayId) - database.outboxDao().deleteGateway(gatewayId) + // The outbox owns command/attachment cascade deletes; nested transactions join this one. + RoomChatCommandOutbox(database).clearGateway(gatewayId) } } } finally { diff --git a/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatCommandOutbox.kt b/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatCommandOutbox.kt index eb395945cd01..6d24e88316e9 100644 --- a/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatCommandOutbox.kt +++ b/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatCommandOutbox.kt @@ -1,7 +1,9 @@ package ai.openclaw.app.chat +import androidx.room.ColumnInfo import androidx.room.Dao import androidx.room.Entity +import androidx.room.Index import androidx.room.Insert import androidx.room.PrimaryKey import androidx.room.Query @@ -20,21 +22,56 @@ internal const val OUTBOX_EXPIRED_ERROR = "expired" /** Delivery is ambiguous after dispatch without an acknowledgement; retry needs explicit intent. */ internal const val OUTBOX_DELIVERY_UNCONFIRMED_ERROR = "delivery unconfirmed; retry manually" +/** Connection-gated command rows never auto-replay across a reconnect; retry needs explicit intent. */ +internal const val OUTBOX_CONNECTION_CHANGED_ERROR = "connection changed before this command was sent; retry manually" + +/** + * gatedEpoch sentinel for rows migrated from schemas without epochs: it matches no live + * connection generation, so legacy command-shaped rows park instead of auto-replaying. + */ +internal const val OUTBOX_GATED_EPOCH_NEVER = -1L + +/** Chunk size for attachment BLOBs; each chunk row must stay well under Android's CursorWindow cap. */ +internal const val OUTBOX_ATTACHMENT_CHUNK_BYTES = 512 * 1024 + +/** Upper bound of attachment bytes on one queued command (8 images plus a voice note fit). */ +internal const val OUTBOX_MAX_COMMAND_ATTACHMENT_BYTES = 8L * 1024L * 1024L + +/** Upper bound of queued attachment bytes per gateway so the outbox database stays bounded. */ +internal const val OUTBOX_MAX_GATEWAY_ATTACHMENT_BYTES = 48L * 1024L * 1024L + enum class ChatOutboxStatus( internal val dbValue: String, ) { Queued("queued"), Sending("sending"), + + /** + * The gateway acknowledged the send, but only canonical chat.history proves the user turn was + * durably persisted (the started ACK is emitted before the transcript write). Accepted rows are + * retired exclusively by history confirmation, or parked as failed when confirmation never lands. + */ + Accepted("accepted"), Failed("failed"), ; internal companion object { - // Destructive migration keeps the schema in lockstep, so unknown values should not occur; - // park anything unexpected as Failed so it stays visible instead of silently sending. + // Schema bumps migrate explicitly, so unknown values should not occur; park anything + // unexpected as Failed so it stays visible instead of silently sending. fun fromDb(value: String): ChatOutboxStatus = entries.firstOrNull { it.dbValue == value } ?: Failed } } +/** Metadata for one durable attachment; bytes live in chunked BLOB rows keyed by [id]. */ +data class ChatOutboxAttachment( + val id: String, + val type: String, + val mimeType: String, + val fileName: String, + val durationMs: Long?, + val byteLength: Long, +) + /** One durable queued chat command; [id] doubles as the chat.send idempotency key. */ data class ChatOutboxItem( val id: String, @@ -47,6 +84,25 @@ data class ChatOutboxItem( val status: ChatOutboxStatus, val retryCount: Int, val lastError: String?, + // Non-null marks a connection-gated row (slash command): it may only auto-send while this + // connection epoch is still active, so reconnects never silently replay a command. + val gatedEpoch: Long? = null, + val attachments: List = emptyList(), +) + +/** Attachment bytes captured at enqueue time; stored as binary chunks, never base64 at rest. */ +class OutboxAttachmentPayload( + val type: String, + val mimeType: String, + val fileName: String, + val durationMs: Long?, + val bytes: ByteArray, +) + +/** One attachment re-assembled for a flush dispatch or a restored optimistic bubble. */ +class LoadedOutboxAttachment( + val attachment: ChatOutboxAttachment, + val bytes: ByteArray, ) sealed interface ChatOutboxEnqueueResult { @@ -56,20 +112,28 @@ sealed interface ChatOutboxEnqueueResult { data object QueueFull : ChatOutboxEnqueueResult + /** One command's attachments exceed [OUTBOX_MAX_COMMAND_ATTACHMENT_BYTES]; deleting rows cannot help. */ + data object AttachmentsTooLarge : ChatOutboxEnqueueResult + + /** The per-gateway attachment byte budget is exhausted; deleting queued rows frees space. */ + data object StorageFull : ChatOutboxEnqueueResult + /** No gateway identity is available (nothing paired/configured), so nothing can be queued. */ data object Unavailable : ChatOutboxEnqueueResult } /** - * Durable offline outbox for text chat commands. + * Durable outbox for chat sends. Every send is journaled here before any network attempt so + * process death always has exactly one recovery owner; rows survive until canonical chat.history + * proves the user turn persisted, they terminally fail, expire, or the user deletes them. * * Unlike the disposable transcript cache, queued rows are user input that must survive process - * restarts until they are acked by the gateway, expired, or explicitly deleted. Like the cache, - * callers bind every gateway-scoped operation to an explicit [ChatCacheScope] gateway id captured - * before their suspend point, so a connection switch cannot re-scope rows mid-operation. + * restarts and schema migrations. Like the cache, callers bind every gateway-scoped operation to + * an explicit [ChatCacheScope] gateway id captured before their suspend point, so a connection + * switch cannot re-scope rows mid-operation. */ interface ChatCommandOutbox { - /** All rows for [gatewayId], strictly createdAt-ordered. */ + /** All rows for [gatewayId] with attachment metadata, strictly createdAt-ordered. */ suspend fun load(gatewayId: String): List suspend fun enqueue( @@ -78,8 +142,13 @@ interface ChatCommandOutbox { text: String, thinkingLevel: String, nowMs: Long, + attachments: List = emptyList(), + gatedEpoch: Long? = null, ): ChatOutboxEnqueueResult + /** Re-assembles the attachment bytes for one command, in stable position order. */ + suspend fun loadAttachments(id: String): List + /** Returns the number of rows updated (0 when the row no longer exists), so callers can claim. */ suspend fun updateStatus( id: String, @@ -88,19 +157,47 @@ interface ChatCommandOutbox { lastError: String?, ): Int + /** + * Atomically claims a queued row for one dispatch (queued -> sending). Returns 0 when the row + * vanished or another dispatcher already claimed it, so the direct-send path and the flush + * loop can never both send the same row. + */ + suspend fun claimForSending( + id: String, + retryCount: Int, + lastError: String?, + ): Int + + /** + * Pins a row enqueued under the pre-hello "main" alias to the canonical session key it first + * resolves to. Replay after that must never re-resolve, so a later default-agent change + * cannot redirect already-captured input. + */ + suspend fun pinSessionKey( + id: String, + sessionKey: String, + ) + /** * User-driven retry of a failed row owned by [gatewayId]: back to 'queued' with reset attempts * and a fresh createdAt, so an expired row is not immediately re-expired by the flush sweep. * Returns the number of rows transitioned; keeps the row id as the gateway idempotency key. + * Gated rows are re-stamped with the caller's current connection epoch, and queued successors + * in the same session shift behind the retried row in their original order, so retrying an + * ambiguous head can never make younger turns of the conversation overtake it. */ suspend fun requeueForRetry( gatewayId: String, id: String, nowMs: Long, + gatedEpoch: Long?, ): Int suspend fun delete(id: String) + /** Retires rows proven delivered by canonical history; returns how many rows were removed. */ + suspend fun confirmDelivered(ids: Set): Int + /** Drops queued commands for a deleted session so they cannot send into a dead session. */ suspend fun deleteForSession( gatewayId: String, @@ -113,7 +210,11 @@ interface ChatCommandOutbox { /** Crash safety: rows stuck in 'sending' after a killed process become visible failed rows. */ suspend fun failSendingAfterRestart() - /** Expires queued rows older than [OUTBOX_EXPIRY_MS] to 'failed' instead of sending stale commands. */ + /** + * Expires stale rows to 'failed' instead of sending stale commands: queued rows older than + * [OUTBOX_EXPIRY_MS] expire, and accepted rows never confirmed within the same window park + * as delivery-unconfirmed so they stay visible for manual review. + */ suspend fun expireStale( gatewayId: String, nowMs: Long, @@ -131,6 +232,32 @@ internal data class OutboxCommandEntity( val status: String, val retryCount: Int, val lastError: String?, + val gatedEpoch: Long?, +) + +@Entity( + tableName = "outbox_attachments", + indices = [Index("commandId")], +) +internal data class OutboxAttachmentEntity( + @PrimaryKey val id: String, + val commandId: String, + val position: Int, + val type: String, + val mimeType: String, + val fileName: String, + val durationMs: Long?, + val byteLength: Long, +) + +@Entity( + tableName = "outbox_attachment_chunks", + primaryKeys = ["attachmentId", "chunkIndex"], +) +internal class OutboxAttachmentChunkEntity( + val attachmentId: String, + val chunkIndex: Int, + @ColumnInfo(typeAffinity = ColumnInfo.BLOB) val bytes: ByteArray, ) @Dao @@ -139,6 +266,15 @@ internal interface ChatOutboxDao { @Query("SELECT * FROM outbox_commands WHERE gatewayId = :gatewayId ORDER BY createdAtMs ASC, id ASC") suspend fun commands(gatewayId: String): List + @Query("SELECT id FROM outbox_commands WHERE gatewayId = :gatewayId AND sessionKey = :sessionKey") + suspend fun commandIdsForSession( + gatewayId: String, + sessionKey: String, + ): List + + @Query("SELECT id FROM outbox_commands WHERE gatewayId = :gatewayId") + suspend fun commandIdsForGateway(gatewayId: String): List + @Query("SELECT COUNT(*) FROM outbox_commands WHERE gatewayId = :gatewayId") suspend fun count(gatewayId: String): Int @@ -156,6 +292,30 @@ internal interface ChatOutboxDao { lastError: String?, ): Int + @Query( + "UPDATE outbox_commands SET status = :toStatus, retryCount = :retryCount, lastError = :lastError " + + "WHERE id = :id AND status = :fromStatus", + ) + suspend fun claimStatus( + id: String, + fromStatus: String, + toStatus: String, + retryCount: Int, + lastError: String?, + ): Int + + @Query("UPDATE outbox_commands SET sessionKey = :sessionKey WHERE id = :id") + suspend fun updateSessionKey( + id: String, + sessionKey: String, + ) + + @Query("UPDATE outbox_commands SET createdAtMs = :createdAtMs WHERE id = :id") + suspend fun updateCreatedAt( + id: String, + createdAtMs: Long, + ) + @Query("UPDATE outbox_commands SET status = :failedStatus, lastError = :error WHERE status = :sendingStatus") suspend fun failAllSending( sendingStatus: String, @@ -164,8 +324,8 @@ internal interface ChatOutboxDao { ) @Query( - "UPDATE outbox_commands SET status = :queuedStatus, retryCount = 0, lastError = NULL, createdAtMs = :createdAtMs " + - "WHERE id = :id AND gatewayId = :gatewayId AND status = :failedStatus", + "UPDATE outbox_commands SET status = :queuedStatus, retryCount = 0, lastError = NULL, createdAtMs = :createdAtMs, " + + "gatedEpoch = :gatedEpoch WHERE id = :id AND gatewayId = :gatewayId AND status = :failedStatus", ) suspend fun requeueForRetry( id: String, @@ -173,54 +333,71 @@ internal interface ChatOutboxDao { createdAtMs: Long, queuedStatus: String, failedStatus: String, + gatedEpoch: Long?, ): Int @Query( "UPDATE outbox_commands SET status = :failedStatus, lastError = :error " + - "WHERE gatewayId = :gatewayId AND status = :queuedStatus AND createdAtMs <= :cutoffMs", + "WHERE gatewayId = :gatewayId AND status = :fromStatus AND createdAtMs <= :cutoffMs", ) - suspend fun expireQueuedAtOrBefore( + suspend fun expireStatusAtOrBefore( gatewayId: String, cutoffMs: Long, - queuedStatus: String, + fromStatus: String, failedStatus: String, error: String, ) @Query("DELETE FROM outbox_commands WHERE id = :id") - suspend fun delete(id: String) + suspend fun delete(id: String): Int - @Query("DELETE FROM outbox_commands WHERE gatewayId = :gatewayId AND sessionKey = :sessionKey") - suspend fun deleteForSession( - gatewayId: String, - sessionKey: String, + @Query("SELECT * FROM outbox_attachments WHERE commandId IN (:commandIds) ORDER BY position ASC") + suspend fun attachmentsForCommands(commandIds: List): List + + @Query("SELECT * FROM outbox_attachments WHERE commandId = :commandId ORDER BY position ASC") + suspend fun attachmentsForCommand(commandId: String): List + + @Query("SELECT bytes FROM outbox_attachment_chunks WHERE attachmentId = :attachmentId ORDER BY chunkIndex ASC") + suspend fun chunksForAttachment(attachmentId: String): List + + @Insert + suspend fun insertAttachment(row: OutboxAttachmentEntity) + + @Insert + suspend fun insertChunk(row: OutboxAttachmentChunkEntity) + + @Query( + "SELECT COALESCE(SUM(byteLength), 0) FROM outbox_attachments WHERE commandId IN " + + "(SELECT id FROM outbox_commands WHERE gatewayId = :gatewayId)", ) + suspend fun attachmentBytesForGateway(gatewayId: String): Long - @Query("DELETE FROM outbox_commands WHERE gatewayId = :gatewayId") - suspend fun deleteGateway(gatewayId: String) + @Query( + "DELETE FROM outbox_attachment_chunks WHERE attachmentId IN " + + "(SELECT id FROM outbox_attachments WHERE commandId = :commandId)", + ) + suspend fun deleteChunksForCommand(commandId: String) + + @Query("DELETE FROM outbox_attachments WHERE commandId = :commandId") + suspend fun deleteAttachmentsForCommand(commandId: String) } /** * Room-backed [ChatCommandOutbox] sharing the chat cache database. Callers pass the gateway id * captured before their suspend point; a blank identity disables both reads and writes. + * Command rows and their attachment bytes are admitted and retired in single transactions, so + * a crash can never orphan bytes or strand a row without its attachments. */ class RoomChatCommandOutbox internal constructor( private val database: ChatCacheDatabase, ) : ChatCommandOutbox { override suspend fun load(gatewayId: String): List { val gateway = scopedGatewayId(gatewayId) ?: return emptyList() - return database.outboxDao().commands(gateway).map { row -> - ChatOutboxItem( - id = row.id, - sessionKey = row.sessionKey, - text = row.text, - thinkingLevel = row.thinkingLevel, - createdAtMs = row.createdAtMs, - status = ChatOutboxStatus.fromDb(row.status), - retryCount = row.retryCount, - lastError = row.lastError, - ) - } + val dao = database.outboxDao() + val rows = dao.commands(gateway) + if (rows.isEmpty()) return emptyList() + val attachmentsByCommand = dao.attachmentsForCommands(rows.map { it.id }).groupBy { it.commandId } + return rows.map { row -> row.toItem(attachmentsByCommand[row.id].orEmpty()) } } override suspend fun enqueue( @@ -229,48 +406,92 @@ class RoomChatCommandOutbox internal constructor( text: String, thinkingLevel: String, nowMs: Long, + attachments: List, + gatedEpoch: Long?, ): ChatOutboxEnqueueResult { val gateway = scopedGatewayId(gatewayId) ?: return ChatOutboxEnqueueResult.Unavailable val key = sessionKey.trim().takeIf { it.isNotEmpty() } ?: return ChatOutboxEnqueueResult.Unavailable + val attachmentBytes = attachments.sumOf { it.bytes.size.toLong() } + if (attachmentBytes > OUTBOX_MAX_COMMAND_ATTACHMENT_BYTES) { + return ChatOutboxEnqueueResult.AttachmentsTooLarge + } val dao = database.outboxDao() - // The bound counts every row (failed included) so total storage stays capped; failed rows - // are user-visible and deletable, so a full queue is always recoverable from the UI. - val row = - database.withTransaction { - if (dao.count(gateway) >= OUTBOX_MAX_QUEUED) { - null - } else { - // Monotonic per-gateway createdAt keeps flush strictly FIFO even when two sends land - // in the same wall-clock millisecond (the id tiebreak is a random UUID otherwise). - val createdAt = maxOf(nowMs, (dao.maxCreatedAt(gateway) ?: Long.MIN_VALUE) + 1) - val entity = - OutboxCommandEntity( + // Admission is one transaction: capacity checks plus the command, attachment, and chunk + // rows commit atomically, so durable admission is all-or-nothing across a crash. The row + // bound counts every row (failed included) so total storage stays capped; failed rows are + // user-visible and deletable, so a full queue is always recoverable from the UI. + return database.withTransaction { + if (dao.count(gateway) >= OUTBOX_MAX_QUEUED) { + return@withTransaction ChatOutboxEnqueueResult.QueueFull + } + if (attachmentBytes > 0 && + dao.attachmentBytesForGateway(gateway) + attachmentBytes > OUTBOX_MAX_GATEWAY_ATTACHMENT_BYTES + ) { + return@withTransaction ChatOutboxEnqueueResult.StorageFull + } + // Monotonic per-gateway createdAt keeps flush strictly FIFO even when two sends land + // in the same wall-clock millisecond (the id tiebreak is a random UUID otherwise). + val createdAt = maxOf(nowMs, (dao.maxCreatedAt(gateway) ?: Long.MIN_VALUE) + 1) + val entity = + OutboxCommandEntity( + id = UUID.randomUUID().toString(), + gatewayId = gateway, + sessionKey = key, + text = text, + thinkingLevel = thinkingLevel, + createdAtMs = createdAt, + status = ChatOutboxStatus.Queued.dbValue, + retryCount = 0, + lastError = null, + gatedEpoch = gatedEpoch, + ) + dao.insert(entity) + val storedAttachments = + attachments.mapIndexed { position, payload -> + val attachmentEntity = + OutboxAttachmentEntity( id = UUID.randomUUID().toString(), - gatewayId = gateway, - sessionKey = key, - text = text, - thinkingLevel = thinkingLevel, - createdAtMs = createdAt, - status = ChatOutboxStatus.Queued.dbValue, - retryCount = 0, - lastError = null, + commandId = entity.id, + position = position, + type = payload.type, + mimeType = payload.mimeType, + fileName = payload.fileName, + durationMs = payload.durationMs, + byteLength = payload.bytes.size.toLong(), ) - dao.insert(entity) - entity + dao.insertAttachment(attachmentEntity) + var chunkIndex = 0 + var offset = 0 + while (offset < payload.bytes.size) { + val end = minOf(offset + OUTBOX_ATTACHMENT_CHUNK_BYTES, payload.bytes.size) + dao.insertChunk( + OutboxAttachmentChunkEntity( + attachmentId = attachmentEntity.id, + chunkIndex = chunkIndex, + bytes = payload.bytes.copyOfRange(offset, end), + ), + ) + chunkIndex += 1 + offset = end + } + attachmentEntity } - } ?: return ChatOutboxEnqueueResult.QueueFull - return ChatOutboxEnqueueResult.Queued( - ChatOutboxItem( - id = row.id, - sessionKey = row.sessionKey, - text = row.text, - thinkingLevel = row.thinkingLevel, - createdAtMs = row.createdAtMs, - status = ChatOutboxStatus.Queued, - retryCount = 0, - lastError = null, - ), - ) + ChatOutboxEnqueueResult.Queued(entity.toItem(storedAttachments)) + } + } + + override suspend fun loadAttachments(id: String): List { + val dao = database.outboxDao() + return dao.attachmentsForCommand(id).map { row -> + val chunks = dao.chunksForAttachment(row.id) + val bytes = ByteArray(chunks.sumOf { it.size }) + var offset = 0 + for (chunk in chunks) { + chunk.copyInto(bytes, offset) + offset += chunk.size + } + LoadedOutboxAttachment(attachment = row.toAttachment(), bytes = bytes) + } } override suspend fun updateStatus( @@ -280,28 +501,84 @@ class RoomChatCommandOutbox internal constructor( lastError: String?, ): Int = database.outboxDao().updateStatus(id = id, status = status.dbValue, retryCount = retryCount, lastError = lastError) + override suspend fun claimForSending( + id: String, + retryCount: Int, + lastError: String?, + ): Int = + database.outboxDao().claimStatus( + id = id, + fromStatus = ChatOutboxStatus.Queued.dbValue, + toStatus = ChatOutboxStatus.Sending.dbValue, + retryCount = retryCount, + lastError = lastError, + ) + + override suspend fun pinSessionKey( + id: String, + sessionKey: String, + ) { + val key = sessionKey.trim().takeIf { it.isNotEmpty() } ?: return + database.outboxDao().updateSessionKey(id = id, sessionKey = key) + } + override suspend fun requeueForRetry( gatewayId: String, id: String, nowMs: Long, + gatedEpoch: Long?, ): Int { val gateway = scopedGatewayId(gatewayId) ?: return 0 val dao = database.outboxDao() return database.withTransaction { - // Same monotonic clamp as enqueue: a retried row re-joins the end of the FIFO queue. - val createdAt = maxOf(nowMs, (dao.maxCreatedAt(gateway) ?: Long.MIN_VALUE) + 1) - dao.requeueForRetry( - id = id, - gatewayId = gateway, - createdAtMs = createdAt, - queuedStatus = ChatOutboxStatus.Queued.dbValue, - failedStatus = ChatOutboxStatus.Failed.dbValue, - ) + val rows = dao.commands(gateway) + val target = rows.firstOrNull { it.id == id } ?: return@withTransaction 0 + // Same monotonic clamp as enqueue: the fresh createdAt keeps the expiry sweep from + // immediately re-failing a retried stale row. + var createdAt = maxOf(nowMs, (dao.maxCreatedAt(gateway) ?: Long.MIN_VALUE) + 1) + val transitioned = + dao.requeueForRetry( + id = id, + gatewayId = gateway, + createdAtMs = createdAt, + queuedStatus = ChatOutboxStatus.Queued.dbValue, + failedStatus = ChatOutboxStatus.Failed.dbValue, + gatedEpoch = gatedEpoch, + ) + if (transitioned > 0) { + // Queued same-session successors follow the retried row in their original order, so + // retrying an ambiguous head cannot let younger conversation turns overtake it. + for (successor in rows) { + val follows = + successor.id != id && + successor.sessionKey == target.sessionKey && + successor.createdAtMs > target.createdAtMs && + ChatOutboxStatus.fromDb(successor.status) == ChatOutboxStatus.Queued + if (follows) { + createdAt += 1 + dao.updateCreatedAt(id = successor.id, createdAtMs = createdAt) + } + } + } + transitioned } } override suspend fun delete(id: String) { - database.outboxDao().delete(id) + database.withTransaction { + deleteCommandRowLocked(id) + } + } + + override suspend fun confirmDelivered(ids: Set): Int { + if (ids.isEmpty()) return 0 + return database.withTransaction { + var removed = 0 + for (id in ids) { + removed += deleteCommandRowLocked(id) + } + removed + } } override suspend fun deleteForSession( @@ -310,12 +587,22 @@ class RoomChatCommandOutbox internal constructor( ) { val gateway = scopedGatewayId(gatewayId) ?: return val key = sessionKey.trim().takeIf { it.isNotEmpty() } ?: return - database.outboxDao().deleteForSession(gateway, key) + val dao = database.outboxDao() + database.withTransaction { + for (id in dao.commandIdsForSession(gateway, key)) { + deleteCommandRowLocked(id) + } + } } override suspend fun clearGateway(gatewayId: String) { val gateway = scopedGatewayId(gatewayId) ?: return - database.outboxDao().deleteGateway(gateway) + val dao = database.outboxDao() + database.withTransaction { + for (id in dao.commandIdsForGateway(gateway)) { + deleteCommandRowLocked(id) + } + } } override suspend fun failSendingAfterRestart() { @@ -333,14 +620,60 @@ class RoomChatCommandOutbox internal constructor( nowMs: Long, ) { val gateway = scopedGatewayId(gatewayId) ?: return - database.outboxDao().expireQueuedAtOrBefore( - gatewayId = gateway, - cutoffMs = nowMs - OUTBOX_EXPIRY_MS, - queuedStatus = ChatOutboxStatus.Queued.dbValue, - failedStatus = ChatOutboxStatus.Failed.dbValue, - error = OUTBOX_EXPIRED_ERROR, - ) + val dao = database.outboxDao() + val cutoff = nowMs - OUTBOX_EXPIRY_MS + database.withTransaction { + dao.expireStatusAtOrBefore( + gatewayId = gateway, + cutoffMs = cutoff, + fromStatus = ChatOutboxStatus.Queued.dbValue, + failedStatus = ChatOutboxStatus.Failed.dbValue, + error = OUTBOX_EXPIRED_ERROR, + ) + // Accepted rows the gateway never confirmed within the window stay visible as failed + // instead of silently occupying the queue forever. + dao.expireStatusAtOrBefore( + gatewayId = gateway, + cutoffMs = cutoff, + fromStatus = ChatOutboxStatus.Accepted.dbValue, + failedStatus = ChatOutboxStatus.Failed.dbValue, + error = OUTBOX_DELIVERY_UNCONFIRMED_ERROR, + ) + } + } + + // Attachment chunk and metadata rows must die with their command row in the same + // transaction; callers wrap this in database.withTransaction. + private suspend fun deleteCommandRowLocked(id: String): Int { + val dao = database.outboxDao() + dao.deleteChunksForCommand(id) + dao.deleteAttachmentsForCommand(id) + return dao.delete(id) } private fun scopedGatewayId(gatewayId: String): String? = gatewayId.trim().takeIf { it.isNotEmpty() } } + +private fun OutboxCommandEntity.toItem(attachments: List): ChatOutboxItem = + ChatOutboxItem( + id = id, + sessionKey = sessionKey, + text = text, + thinkingLevel = thinkingLevel, + createdAtMs = createdAtMs, + status = ChatOutboxStatus.fromDb(status), + retryCount = retryCount, + lastError = lastError, + gatedEpoch = gatedEpoch, + attachments = attachments.map { it.toAttachment() }, + ) + +private fun OutboxAttachmentEntity.toAttachment(): ChatOutboxAttachment = + ChatOutboxAttachment( + id = id, + type = type, + mimeType = mimeType, + fileName = fileName, + durationMs = durationMs, + byteLength = byteLength, + ) diff --git a/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatController.kt b/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatController.kt index 2dcef42faf5d..3e715e2f63f9 100644 --- a/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatController.kt +++ b/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatController.kt @@ -14,6 +14,7 @@ import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CoroutineStart import kotlinx.coroutines.Job +import kotlinx.coroutines.async import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow @@ -28,6 +29,7 @@ import kotlinx.serialization.json.JsonNull import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonPrimitive import kotlinx.serialization.json.buildJsonObject +import java.util.Base64 import java.util.Locale import java.util.UUID import java.util.concurrent.ConcurrentHashMap @@ -200,6 +202,17 @@ class ChatController internal constructor( private val outboxRecoveryMutex = Mutex() private var outboxRecoveryComplete = false + // Counts idle-history snapshots that lacked proof for an orphaned accepted row; rows park as + // delivery-unconfirmed on the second sighting so one lagging transcript write is not loss. + private val unconfirmedSightings = ConcurrentHashMap() + + // Gateway ACKs may return a run id that differs from the row's idempotency key; ownership + // and in-flight checks must recognize both or reconciliation can park a still-live run. + // Deliberately in-memory: chat.send uses the client idempotency key as the run id, and + // after a restart canonical-history proof by ":user" retires rows regardless of the + // acked id; an ambiguous survivor parks for manual review instead of auto-retrying. + private val acknowledgedRunIdByRowId = ConcurrentHashMap() + private val outboxRecoveryJob = commandOutbox?.let { outbox -> scope.launch { @@ -240,6 +253,9 @@ class ChatController internal constructor( _streamingAssistantText.value = null _historyLoading.value = false _sessionId.value = null + // Failed connect attempts pass through onGatewayScopeChanging, which empties the published + // outbox rows; repopulate for the still-selected gateway so queued sends stay visible offline. + scope.launch { publishOutbox() } } /** Refreshes the connected gateway while preserving recovery ownership after a disconnect. */ @@ -759,7 +775,7 @@ class ChatController internal constructor( } } - /** Sends a chat message and returns once the gateway accepts or rejects the request. */ + /** Sends a chat message and returns once it is durably admitted or the gateway rejects it. */ suspend fun sendMessageAwaitAcceptance( message: String, thinkingLevel: String, @@ -784,43 +800,80 @@ class ChatController internal constructor( } else { "off" } - if (!_healthOk.value) { - // Offline capture: text-only commands become durable outbox rows and flush on reconnect. - // Attachments stay blocked (text-only v1) so large payloads never sit in the database. - if (commandOutbox == null || attachments.isNotEmpty()) { - updateErrorText("Gateway health not OK; cannot send") - return false - } - return enqueueOfflineCommand(text = trimmed, thinkingLevel = thinking) - } - - val runId = UUID.randomUUID().toString() val text = if (trimmed.isEmpty() && attachments.isNotEmpty()) "See attached." else trimmed - // Optimistic user message keeps the composer responsive while chat.send and history refresh complete. - val userContent = - buildList { - add(ChatMessageContent(type = "text", text = text)) - for (att in attachments) { - add( - ChatMessageContent( - type = att.type, - mimeType = att.mimeType, - fileName = att.fileName, - base64 = att.base64, - durationMs = att.durationMs, - ), - ) + // Every send is journaled before the composer clears or any network attempt can lose + // ownership; the durable row is the single recovery owner across process death. + val journaled = + when (val outbox = commandOutbox) { + null -> { + if (!_healthOk.value) { + updateErrorText("Gateway health not OK; cannot send") + return false + } + null } + else -> + enqueueDurableSend( + outbox = outbox, + outboxScope = sendCacheScope, + sessionKey = normalizeRequestedSessionKey(sessionKey), + text = text, + thinkingLevel = thinking, + attachments = attachments, + ) ?: return false } - val optimisticMessage = - ChatMessage( - id = UUID.randomUUID().toString(), - role = "user", - content = userContent, - timestampMs = System.currentTimeMillis(), - idempotencyKey = "$runId:user", - ) + if (journaled != null) { + if (!_healthOk.value) { + // Captured for reconnect: the queued bubble is visible and flush delivers it later. + return true + } + // The startup recovery sweep flips every 'sending' row to delivery-unconfirmed. Claiming + // only after it completes means the sweep can never hit this live dispatch; a failed + // sweep leaves the row queued so reconnect flush owns delivery instead. + outboxRecoveryJob?.join() + val outbox = commandOutbox + if (outbox == null || !recoverInterruptedOutboxSends(outbox)) { + _healthOk.value = false + publishOutbox() + return true + } + if (sessionHasDurableBacklog(journaled)) { + // An older row for this session is still queued or unresolved; a direct dispatch + // would reorder the conversation, so the FIFO flush owns delivery. + requestOutboxFlush() + return true + } + // Atomically claim the row for this direct dispatch: a vanished row (user delete) or a + // concurrent flush claim must not lead to a second send of the same idempotency key. + val claimed = + try { + outbox.claimForSending(journaled.id, 0, null) + } catch (err: CancellationException) { + throw err + } catch (_: Throwable) { + null + } + publishOutbox() + if (claimed == null) { + // The claim could not be made durable, so the admitted row still has no dispatcher. + // Hand delivery to the flush lane instead of reporting success with no active owner. + requestOutboxFlush() + return true + } + if (claimed == 0) return true + if (journaled.gatedEpoch != null && journaled.gatedEpoch != currentCacheScope()?.connectionGeneration) { + // A reconnect landed between admission and this claim; command-shaped input never + // auto-replays across connection epochs, so the claimed row parks for explicit retry. + persistJournaledSendState(journaled, ChatOutboxStatus.Failed, OUTBOX_CONNECTION_CHANGED_ERROR) + return true + } + } + + val runId = journaled?.id ?: UUID.randomUUID().toString() + + // Optimistic user message keeps the composer responsive while chat.send and history refresh complete. + val optimisticMessage = optimisticUserMessage(runId = runId, text = text, attachments = attachments) optimisticMessagesByRunId[runId] = optimisticMessage unresolvedRepliesByRunId[runId] = optimisticMessage _messages.value = _messages.value + optimisticMessage @@ -836,85 +889,263 @@ class ChatController internal constructor( pendingToolCallsById.clear() publishPendingToolCalls() - return try { - val params = - buildJsonObject { - put("sessionKey", JsonPrimitive(sessionKey)) - put("message", JsonPrimitive(text)) - put("thinking", JsonPrimitive(thinking)) - put("timeoutMs", JsonPrimitive(30_000)) - put("idempotencyKey", JsonPrimitive(runId)) - if (attachments.isNotEmpty()) { - put( - "attachments", - JsonArray( - attachments.map { att -> - buildJsonObject { - put("type", JsonPrimitive(att.type)) - put("mimeType", JsonPrimitive(att.mimeType)) - put("fileName", JsonPrimitive(att.fileName)) - put("content", JsonPrimitive(att.base64)) - } - }, - ), + // Dispatch ownership lives in the controller scope: cancelling the calling UI scope + // (leaving the chat screen mid-send) after the durable claim must not strand a Sending + // row this process can no longer repair; the dispatch completes and settles the row. + val dispatch = + scope.async { + try { + val params = + buildChatSendParams( + // Dispatch exactly what was journaled: the row's captured session key is the + // idempotent identity a replay after process death would use. + sessionKey = journaled?.sessionKey ?: sessionKey, + text = text, + thinking = thinking, + idempotencyKey = runId, + attachments = attachments, ) + val res = requestGatewayBound(sendGatewayId, "chat.send", params) + val ack = parseChatSendAck(json, res) + // Row transitions are durable state for the dispatching gateway and apply even when the + // UI scope moved on mid-request; only UI updates below are scope-guarded. A terminal + // failure ack proves transmission, not that this idempotency key never ran (a timeout ack + // can outlive a still-admitted run), so the row parks for review instead of deleting. + if (ack.isTerminalFailure) { + markJournaledSendUnconfirmed(journaled) + } else { + markJournaledSendAccepted(journaled) + val ackRunId = ack.runId + if (journaled != null && ackRunId != null && ackRunId != journaled.id) { + acknowledgedRunIdByRowId[journaled.id] = ackRunId + } + } + if (sendCacheScope != currentCacheScope()) return@async true + val actualRunId = ack.runId ?: runId + if (actualRunId != runId) { + transferRunOwnership(runId, actualRunId, optimisticMessage) + } + if (ack.isTerminal) { + clearPendingRun(actualRunId) + removeOptimisticMessage(actualRunId) + pendingToolCallsById.clear() + publishPendingToolCalls() + _streamingAssistantText.value = null + if (ack.isTerminalSuccess) { + unresolvedRepliesByRunId.remove(actualRunId) + refreshCurrentHistoryBestEffort(runIdsToReconcile = setOf(actualRunId)) + true + } else { + // Terminal timeout/error means the gateway did not accept a runnable turn. + // Surface failed acceptance instead of letting a cleared composer look successful. + unresolvedRepliesByRunId.remove(actualRunId) + updateErrorText("Chat failed before the run started; try again.") + // The parked row owns the input; restoring the draft would duplicate it. + journaled != null + } + } else { + true + } + } catch (err: CancellationException) { + throw err + } catch (err: GatewayRequestNotEnqueued) { + // The frame provably never entered the socket queue. The journaled row stays queued and + // reconnect flush owns delivery, exactly like the flush path treats not-dispatched sends; + // deleting here could lose fire-and-forget input if the process died after the delete. + if (journaled != null) { + persistJournaledSendState(journaled, ChatOutboxStatus.Queued, err.message) + if (sendCacheScope != currentCacheScope()) return@async true + clearPendingRun(runId) + removeOptimisticMessage(runId) + unresolvedRepliesByRunId.remove(runId) + // The transport is effectively down; drop health so the next health event re-flushes. + _healthOk.value = false + publishOutbox() + true + } else { + if (sendCacheScope != currentCacheScope()) return@async true + clearPendingRun(runId) + removeOptimisticMessage(runId) + unresolvedRepliesByRunId.remove(runId) + updateErrorText(err.message) + false + } + } catch (err: GatewayRequestDefinitiveFailure) { + // An ok:false response proves transmission, not that this idempotency key was never run; + // park the journaled copy for review instead of deleting a possibly delivered send. + markJournaledSendUnconfirmed(journaled) + if (sendCacheScope != currentCacheScope()) return@async true + clearPendingRun(runId) + removeOptimisticMessage(runId) + unresolvedRepliesByRunId.remove(runId) + updateErrorText(err.message) + // The parked row owns the input; only the journal-less path refuses the send. + journaled != null + } catch (_: GatewayRequestOutcomeUnknown) { + // A transport failure cannot distinguish rejection from an accepted send whose ACK was + // lost. Keep the journaled row until history confirms or reconciliation parks it. + markJournaledSendAccepted(journaled) + if (sendCacheScope != currentCacheScope()) return@async true + unknownOutcomeRunIds.add(runId) + if (_healthOk.value) { + refreshCurrentHistoryBestEffort(runIdsToReconcile = setOf(runId)) } - } - val res = requestGatewayBound(sendGatewayId, "chat.send", params.toString()) - if (sendCacheScope != currentCacheScope()) return true - val ack = parseChatSendAck(json, res) - val actualRunId = ack.runId ?: runId - if (actualRunId != runId) { - transferRunOwnership(runId, actualRunId, optimisticMessage) - } - if (ack.isTerminal) { - clearPendingRun(actualRunId) - removeOptimisticMessage(actualRunId) - pendingToolCallsById.clear() - publishPendingToolCalls() - _streamingAssistantText.value = null - if (ack.isTerminalSuccess) { - unresolvedRepliesByRunId.remove(actualRunId) - refreshCurrentHistoryBestEffort(runIdsToReconcile = setOf(actualRunId)) true - } else { - // Terminal timeout/error means the gateway did not accept a runnable turn. - // Surface failed acceptance instead of letting a cleared composer look successful. - unresolvedRepliesByRunId.remove(actualRunId) - updateErrorText("Chat failed before the run started; try again.") - false + } catch (err: Throwable) { + // Unexpected failure after dispatch is ambiguous; fail closed and keep the row visible. + markJournaledSendUnconfirmed(journaled) + if (sendCacheScope != currentCacheScope()) return@async true + clearPendingRun(runId) + removeOptimisticMessage(runId) + unresolvedRepliesByRunId.remove(runId) + updateErrorText(err.message) + // With a journaled row parked for review, the composer must not restore a duplicate + // draft: the row owns the input now. Only the journal-less path refuses the send. + journaled != null } - } else { - true } - } catch (err: CancellationException) { - throw err - } catch (err: GatewayRequestDefinitiveFailure) { - if (sendCacheScope != currentCacheScope()) return true - clearPendingRun(runId) - removeOptimisticMessage(runId) - unresolvedRepliesByRunId.remove(runId) - updateErrorText(err.message) - false - } catch (_: GatewayRequestOutcomeUnknown) { - if (sendCacheScope != currentCacheScope()) return true - // A transport failure cannot distinguish rejection from an accepted send whose - // ACK was lost. Keep the idempotency-key-backed row to prevent a duplicate retry. - unknownOutcomeRunIds.add(runId) - if (_healthOk.value) { - refreshCurrentHistoryBestEffort(runIdsToReconcile = setOf(runId)) + return dispatch.await() + } + + private fun optimisticUserMessage( + runId: String, + text: String, + attachments: List, + ): ChatMessage { + val userContent = + buildList { + add(ChatMessageContent(type = "text", text = text)) + for (att in attachments) { + add( + ChatMessageContent( + type = att.type, + mimeType = att.mimeType, + fileName = att.fileName, + base64 = att.base64, + durationMs = att.durationMs, + ), + ) + } } - true - } catch (err: Throwable) { - if (sendCacheScope != currentCacheScope()) return true - clearPendingRun(runId) - removeOptimisticMessage(runId) - unresolvedRepliesByRunId.remove(runId) - updateErrorText(err.message) - false + return ChatMessage( + id = UUID.randomUUID().toString(), + role = "user", + content = userContent, + timestampMs = System.currentTimeMillis(), + idempotencyKey = "$runId:user", + ) + } + + private fun buildChatSendParams( + sessionKey: String, + text: String, + thinking: String, + idempotencyKey: String, + attachments: List, + ): String = + buildJsonObject { + put("sessionKey", JsonPrimitive(sessionKey)) + put("message", JsonPrimitive(text)) + put("thinking", JsonPrimitive(thinking)) + put("timeoutMs", JsonPrimitive(30_000)) + put("idempotencyKey", JsonPrimitive(idempotencyKey)) + if (attachments.isNotEmpty()) { + put( + "attachments", + JsonArray( + attachments.map { att -> + buildJsonObject { + put("type", JsonPrimitive(att.type)) + put("mimeType", JsonPrimitive(att.mimeType)) + put("fileName", JsonPrimitive(att.fileName)) + put("content", JsonPrimitive(att.base64)) + } + }, + ), + ) + } + }.toString() + + /** True when an older durable row for the same session must send before this one. */ + private suspend fun sessionHasDurableBacklog(row: ChatOutboxItem): Boolean { + val outbox = commandOutbox ?: return false + val outboxScope = currentCacheScope() ?: return false + val rows = runCatching { outbox.load(outboxScope.gatewayId) }.getOrDefault(emptyList()) + return rows.any { other -> + other.id != row.id && + other.createdAtMs < row.createdAtMs && + sameOutboxSession(other.sessionKey, row.sessionKey) && + outboxRowUnresolved(other) } } + // Queued/sending rows are still ahead in FIFO order, and an orphaned accepted row holds its + // session only until history proof confirms or parks it (a bounded window). Parked failed + // rows are terminal-manual state and do not strand later turns; explicit Retry re-orders + // still-queued successors behind the retried head instead. + private fun outboxRowUnresolved(row: ChatOutboxItem): Boolean = + when (row.status) { + ChatOutboxStatus.Queued, ChatOutboxStatus.Sending -> true + ChatOutboxStatus.Accepted -> !locallyOwnedOutboxRow(row.id) + ChatOutboxStatus.Failed -> false + } + + // A row is live-owned when either its idempotency key or the run id the gateway + // acknowledged it under still has local pending/unknown/unresolved state. + private fun locallyOwnedOutboxRow(rowId: String): Boolean = locallyOwnedRun(rowId) || acknowledgedRunIdByRowId[rowId]?.let(::locallyOwnedRun) == true + + private fun locallyOwnedRun(runId: String): Boolean = + synchronized(pendingRuns) { pendingRuns.contains(runId) } || + unknownOutcomeRunIds.contains(runId) || + unresolvedRepliesByRunId.containsKey(runId) + + private fun sameOutboxSession( + left: String, + right: String, + ): Boolean = normalizeRequestedSessionKey(left) == normalizeRequestedSessionKey(right) + + private suspend fun markJournaledSendAccepted(row: ChatOutboxItem?) { + persistJournaledSendState(row, ChatOutboxStatus.Accepted, null) + } + + private suspend fun markJournaledSendUnconfirmed(row: ChatOutboxItem?) { + persistJournaledSendState(row, ChatOutboxStatus.Failed, OUTBOX_DELIVERY_UNCONFIRMED_ERROR) + } + + // Mirrors the flush path's fail-closed persistence handling: a claimed row whose follow-up + // state cannot be made durable must not silently stay 'sending' (it would block its session + // with no user action available); the re-armed recovery sweep parks it once storage recovers. + private suspend fun persistJournaledSendState( + row: ChatOutboxItem?, + status: ChatOutboxStatus, + lastError: String?, + ) { + val outbox = commandOutbox ?: return + if (row == null) return + if (status != ChatOutboxStatus.Accepted) acknowledgedRunIdByRowId.remove(row.id) + val persisted = + try { + outbox.updateStatus(row.id, status, row.retryCount, lastError) + } catch (err: CancellationException) { + throw err + } catch (_: Throwable) { + null + } + if (persisted == null) { + rearmOutboxRecovery() + _healthOk.value = false + } + publishOutbox() + kickFlushForRoutedBacklog() + } + + // Sends routed to the queue while a direct dispatch held their session wait for that dispatch + // to resolve; re-kick the single-flight flush so they do not idle until the next health event. + private fun kickFlushForRoutedBacklog() { + if (!_healthOk.value) return + requestOutboxFlush() + } + /** Sends best-effort abort requests for every currently pending gateway run. */ fun abort() { val abortGatewayId = currentCacheScope()?.gatewayId @@ -1182,9 +1413,24 @@ class ChatController internal constructor( if (!applied) return false completeReconnectRecoveryIfOwned(sessionKey, generation) persistTranscript(requestCacheScope, sessionKey, history.messages) + confirmDurableSendsFromHistory(requestCacheScope, history) return true } + /** Canonical history is the only proof that retires journaled sends; every apply checks it. */ + private suspend fun confirmDurableSendsFromHistory( + requestCacheScope: ChatCacheScope?, + history: ChatHistory, + ) { + val outbox = commandOutbox ?: return + val gatewayId = requestCacheScope?.gatewayId ?: return + if (reconcileDurableSendsAgainstHistory(outbox, gatewayId, history)) { + publishOutbox() + // Retired rows may have been session heads holding queued successors; resume delivery. + kickFlushForRoutedBacklog() + } + } + /** Lets whichever same-generation history request wins finish reconnect health recovery. */ private suspend fun completeReconnectRecoveryIfOwned( sessionKey: String, @@ -1402,42 +1648,78 @@ class ChatController internal constructor( scope.launch { fetchChatMetadata() } } - private suspend fun enqueueOfflineCommand( + /** + * Durably admits one send (text plus decoded attachment bytes) before any network attempt. + * Returns null after surfacing an actionable error; the composer must keep the draft then. + */ + private suspend fun enqueueDurableSend( + outbox: ChatCommandOutbox, + outboxScope: ChatCacheScope?, + sessionKey: String, text: String, thinkingLevel: String, - ): Boolean { - val outbox = commandOutbox ?: return false - val outboxScope = - currentCacheScope() ?: run { - updateErrorText("Gateway health not OK; cannot send") - return false + attachments: List, + ): ChatOutboxItem? { + if (outboxScope == null) { + updateErrorText("Gateway health not OK; cannot send") + return null + } + val payloads = + try { + attachments.map { att -> + OutboxAttachmentPayload( + type = att.type, + mimeType = att.mimeType, + fileName = att.fileName, + durationMs = att.durationMs, + bytes = Base64.getDecoder().decode(att.base64), + ) + } + } catch (_: IllegalArgumentException) { + updateErrorText("Could not stage an attachment for sending.") + return null } + // Slash commands are connection-gated: they may auto-send only inside the connection epoch + // that captured them, so a reconnect never silently replays a command-shaped input. + val gatedEpoch = if (text.startsWith("/")) outboxScope.connectionGeneration else null val result = try { outbox.enqueue( gatewayId = outboxScope.gatewayId, - sessionKey = _sessionKey.value, + sessionKey = sessionKey, text = text, thinkingLevel = thinkingLevel, nowMs = System.currentTimeMillis(), + attachments = payloads, + gatedEpoch = gatedEpoch, ) + } catch (err: CancellationException) { + throw err } catch (_: Throwable) { updateErrorText("Could not queue message for later delivery.") - return false + return null } return when (result) { is ChatOutboxEnqueueResult.Queued -> { updateErrorText(null) publishOutbox() - true + result.item } ChatOutboxEnqueueResult.QueueFull -> { updateErrorText("Offline queue is full ($OUTBOX_MAX_QUEUED messages); delete queued items first.") - false + null + } + ChatOutboxEnqueueResult.AttachmentsTooLarge -> { + updateErrorText("Attachments are too large to queue for one message; remove some and try again.") + null + } + ChatOutboxEnqueueResult.StorageFull -> { + updateErrorText("Offline attachment storage is full; delete queued items first.") + null } ChatOutboxEnqueueResult.Unavailable -> { updateErrorText("Gateway health not OK; cannot send") - false + null } } } @@ -1447,11 +1729,21 @@ class ChatController internal constructor( val outbox = commandOutbox ?: return scope.launch { val outboxScope = currentCacheScope() ?: return@launch + val row = _outboxItems.value.firstOrNull { it.id == id } + // A gated command row is re-armed for the current connection epoch only; retrying it + // while disconnected parks it again at the next reconnect instead of silently replaying. + val gatedEpoch = row?.gatedEpoch?.let { outboxScope.connectionGeneration } // requeueForRetry refreshes createdAt and requires this gateway's Failed state. The // compare-and-set keeps stale gateway or double Retry taps from reviving an in-flight row. val requeued = - runCatching { outbox.requeueForRetry(gatewayId = outboxScope.gatewayId, id = id, nowMs = System.currentTimeMillis()) } - .getOrDefault(0) + runCatching { + outbox.requeueForRetry( + gatewayId = outboxScope.gatewayId, + id = id, + nowMs = System.currentTimeMillis(), + gatedEpoch = gatedEpoch, + ) + }.getOrDefault(0) publishOutbox() if (requeued > 0 && _healthOk.value) requestOutboxFlush() } @@ -1461,7 +1753,10 @@ class ChatController internal constructor( val outbox = commandOutbox ?: return scope.launch { runCatching { outbox.delete(id) } + acknowledgedRunIdByRowId.remove(id) publishOutbox() + // Deleting an unresolved row can release its session's queued successors. + if (_healthOk.value) requestOutboxFlush() } } @@ -1522,16 +1817,29 @@ class ChatController internal constructor( runCatching { outbox.expireStale(flushScope.gatewayId, System.currentTimeMillis()) } publishOutbox() while (_healthOk.value && currentCacheScope() == flushScope) { - val next = - runCatching { outbox.load(flushScope.gatewayId) } - .getOrDefault(emptyList()) - .firstOrNull { it.status == ChatOutboxStatus.Queued } ?: break + val rows = runCatching { outbox.load(flushScope.gatewayId) }.getOrDefault(emptyList()) + if (parkStaleGatedRows(outbox, rows, flushScope)) { + publishOutbox() + continue + } + val next = nextFlushableRow(rows) ?: break when (sendOutboxItem(outbox, next, flushScope)) { OutboxSendOutcome.Sent -> flushedAny = true OutboxSendOutcome.Continue -> {} OutboxSendOutcome.Stop -> break } } + // Accepted rows from an earlier process have no live run ownership; prove them against + // canonical history now so restarts either retire them or surface them for review. The + // second pass (after a short delay) both confirms turns whose transcript write lagged the + // ACK and provides the second sighting that parks genuinely lost sends. Confirmations can + // release queued successors in the same session, so they request a rerun of the drain. + if (reconcileOrphanAcceptedRows(outbox, flushScope) > 0) { + delay(recoveryHistoryRetryDelayMs) + if (_healthOk.value && currentCacheScope() == flushScope) { + reconcileOrphanAcceptedRows(outbox, flushScope) + } + } } finally { publishOutbox() if (flushedAny) { @@ -1541,6 +1849,163 @@ class ChatController internal constructor( } } + /** + * First queued row whose session has no earlier unresolved row. Rows are createdAt-ordered, so + * an unresolved row (queued behind a dispatch, ambiguous, or awaiting proof) holds only its own + * session while other sessions keep flushing. + */ + private fun nextFlushableRow(rows: List): ChatOutboxItem? { + val blockedSessions = mutableSetOf() + for (row in rows) { + val session = normalizeRequestedSessionKey(row.sessionKey) + if (row.status == ChatOutboxStatus.Queued && session !in blockedSessions) return row + if (outboxRowUnresolved(row)) blockedSessions.add(session) + } + return null + } + + // Gated command rows enqueued under an older connection epoch park instead of auto-replaying; + // returns true when any row changed so the flush loop reloads before selecting. + private suspend fun parkStaleGatedRows( + outbox: ChatCommandOutbox, + rows: List, + flushScope: ChatCacheScope, + ): Boolean { + var parked = false + for (row in rows) { + val stale = + row.status == ChatOutboxStatus.Queued && + row.gatedEpoch != null && + row.gatedEpoch != flushScope.connectionGeneration + if (!stale) continue + // A park that cannot be persisted must fail closed: reporting it as parked would make + // the flush loop reload the same queued row and spin while health stays OK. + val persisted = updateOutboxStatusOrNull(outbox, row, ChatOutboxStatus.Failed, OUTBOX_CONNECTION_CHANGED_ERROR) + if (persisted == null) { + // Returning true here re-enters the loop, whose health check now stops the pass; + // falling through instead would dispatch the still-queued stale row this pass. + rearmOutboxRecovery() + _healthOk.value = false + return true + } + parked = true + } + return parked + } + + /** Reconciles orphaned accepted rows against per-session history; returns how many remain. */ + private suspend fun reconcileOrphanAcceptedRows( + outbox: ChatCommandOutbox, + flushScope: ChatCacheScope, + ): Int { + val rows = runCatching { outbox.load(flushScope.gatewayId) }.getOrDefault(emptyList()) + val orphanSessions = + rows + .filter { it.status == ChatOutboxStatus.Accepted && !locallyOwnedOutboxRow(it.id) } + .map { normalizeRequestedSessionKey(it.sessionKey) } + .toSet() + if (orphanSessions.isEmpty()) return 0 + var changed = false + for (sessionKey in orphanSessions) { + if (!_healthOk.value || currentCacheScope() != flushScope) break + val history = + try { + val historyJson = + requestGatewayBound( + flushScope.gatewayId, + "chat.history", + buildJsonObject { put("sessionKey", JsonPrimitive(sessionKey)) }.toString(), + ) + parseHistory(historyJson, sessionKey = sessionKey, previousMessages = emptyList()) + } catch (err: CancellationException) { + throw err + } catch (_: Throwable) { + // Keep the rows accepted; the next flush or history apply reconciles them. + continue + } + changed = reconcileDurableSendsAgainstHistory(outbox, flushScope.gatewayId, history) || changed + } + if (changed) { + publishOutbox() + // A confirmed row may have been the head blocking queued successors in its session; + // the level-triggered request makes the drain run another pass so released rows send. + outboxFlushRequested.set(true) + } + return runCatching { outbox.load(flushScope.gatewayId) } + .getOrDefault(emptyList()) + .count { it.status == ChatOutboxStatus.Accepted && !locallyOwnedOutboxRow(it.id) } + } + + /** + * Applies canonical history proof to durable rows: any row whose `id:user` idempotency key is + * persisted retires (regardless of state; proof always wins so a manual retry of an actually + * delivered row can never double-send). Orphaned accepted rows absent from an idle history are + * parked as delivery-unconfirmed only after two independent sightings, so a transcript write + * that briefly lags the ACK is not misread as loss. + */ + private suspend fun reconcileDurableSendsAgainstHistory( + outbox: ChatCommandOutbox, + gatewayId: String, + history: ChatHistory, + ): Boolean { + val rows = runCatching { outbox.load(gatewayId) }.getOrDefault(emptyList()) + if (rows.isEmpty()) return false + val provenIds = history.messages.mapNotNull(::outboxRowIdFromMessage).toSet() + val inFlightRunId = + history.inFlightRun + ?.runId + ?.trim() + ?.takeIf { it.isNotEmpty() } + val sessionRows = rows.filter { sameOutboxSession(it.sessionKey, history.sessionKey) } + var changed = false + val confirmed = sessionRows.filter { it.id in provenIds }.map { it.id }.toSet() + if (confirmed.isNotEmpty()) { + val removed = runCatching { outbox.confirmDelivered(confirmed) }.getOrDefault(0) + confirmed.forEach(unconfirmedSightings::remove) + confirmed.forEach(acknowledgedRunIdByRowId::remove) + changed = removed > 0 + } + for (row in sessionRows) { + if (row.status != ChatOutboxStatus.Accepted || row.id in confirmed) continue + if (locallyOwnedOutboxRow(row.id)) continue + // inFlightRunId must be non-null before the map compare: a missing in-flight run would + // otherwise match rows with no acknowledged id (null == null) and block parking forever. + val rowInFlight = + inFlightRunId != null && + (row.id == inFlightRunId || acknowledgedRunIdByRowId[row.id] == inFlightRunId) + if (rowInFlight) { + // The run is still alive on the gateway; its user turn persists with the run. + unconfirmedSightings.remove(row.id) + continue + } + val sightings = (unconfirmedSightings[row.id] ?: 0) + 1 + if (sightings >= 2) { + val persisted = updateOutboxStatusOrNull(outbox, row, ChatOutboxStatus.Failed, OUTBOX_DELIVERY_UNCONFIRMED_ERROR) + if (persisted == null) { + // The park write failed; reporting a change anyway would spin confirm/park passes + // against unavailable storage while the row's session stays blocked. + rearmOutboxRecovery() + _healthOk.value = false + } else { + unconfirmedSightings.remove(row.id) + acknowledgedRunIdByRowId.remove(row.id) + changed = true + } + } else { + unconfirmedSightings[row.id] = sightings + } + } + return changed + } + + /** Extracts the outbox row id from a persisted user turn's `:user` idempotency key. */ + private fun outboxRowIdFromMessage(message: ChatMessage): String? { + if (message.role.trim().lowercase() != "user") return null + val key = message.idempotencyKey?.trim() ?: return null + if (!key.endsWith(":user")) return null + return key.removeSuffix(":user").takeIf { it.isNotEmpty() } + } + // Sent: acked and removed. Continue: row vanished or failed after a gateway response. // Stop: transport or persistence state cannot safely advance to younger work. private enum class OutboxSendOutcome { Sent, Continue, Stop } @@ -1548,7 +2013,9 @@ class ChatController internal constructor( private enum class GatewayResponseState { Received, Unknown } private sealed interface OutboxSendResult { - data object Accepted : OutboxSendResult + data class Accepted( + val runId: String, + ) : OutboxSendResult /** The request never entered the socket queue, so reconnect may retry it automatically. */ data class NotDispatched( @@ -1575,15 +2042,26 @@ class ChatController internal constructor( null } + private suspend fun claimOutboxRowOrNull( + outbox: ChatCommandOutbox, + item: ChatOutboxItem, + ): Int? = + try { + outbox.claimForSending(item.id, item.retryCount, item.lastError) + } catch (err: CancellationException) { + throw err + } catch (_: Throwable) { + null + } + private suspend fun sendOutboxItem( outbox: ChatCommandOutbox, item: ChatOutboxItem, flushScope: ChatCacheScope, ): OutboxSendOutcome { - // Claim the row before sending: 0 updated rows means it was deleted since the load, and a - // deleted command must never be sent. Continue (like an acknowledged failure) lets the - // flush advance to younger rows without replaying this one. - val claimed = updateOutboxStatusOrNull(outbox, item, ChatOutboxStatus.Sending, item.lastError) + // Atomically claim the row before sending: null means the claim could not be made durable, + // and 0 means the row vanished or a direct dispatch claimed it first; neither may dispatch. + val claimed = claimOutboxRowOrNull(outbox, item) publishOutbox() if (claimed == null) { // Never bypass an older row when its claim could not be made durable. @@ -1591,25 +2069,41 @@ class ChatController internal constructor( return OutboxSendOutcome.Stop } if (claimed == 0) return OutboxSendOutcome.Continue - return when (val result = attemptOutboxSend(item, flushScope.gatewayId)) { - OutboxSendResult.Accepted -> { - // Ack received: delete the row so the flushed history copy is the only bubble left. - val deleted = - try { - outbox.delete(item.id) - true - } catch (err: CancellationException) { - throw err - } catch (_: Throwable) { - false - } - if (!deleted) rearmOutboxRecovery() + // Bytes are loaded once per item; a storage failure here parks the row instead of sending + // a message without the attachments the user staged with it. + val attachments = + try { + loadOutboxAttachmentsForSend(outbox, item) + } catch (err: CancellationException) { + throw err + } catch (_: Throwable) { + val parked = updateOutboxStatusOrNull(outbox, item, ChatOutboxStatus.Failed, "attachments unavailable") + if (parked == null) rearmOutboxRecovery() publishOutbox() - if (deleted) { - OutboxSendOutcome.Sent - } else { + return if (parked == null) { _healthOk.value = false OutboxSendOutcome.Stop + } else { + OutboxSendOutcome.Continue + } + } + return when (val result = attemptOutboxSend(outbox, item, flushScope.gatewayId, attachments)) { + is OutboxSendResult.Accepted -> { + // Ack received: keep the row as accepted until canonical history proves the user turn + // persisted; the started ACK alone is not durable proof (issue #86946 tracks the gap). + if (result.runId != item.id) acknowledgedRunIdByRowId[item.id] = result.runId + val persisted = updateOutboxStatusOrNull(outbox, item, ChatOutboxStatus.Accepted, null) + if (persisted == null) rearmOutboxRecovery() + publishOutbox() + if (persisted == null) { + // The accepted row is still Sending; the re-armed recovery sweep parks it once + // storage recovers, and canonical history proof can still retire it later. + _healthOk.value = false + OutboxSendOutcome.Stop + } else { + // A zero update means a concurrent delete raced the ack; history still owns proof. + if (persisted > 0) adoptFlushedSend(item, attachments, result.runId) + OutboxSendOutcome.Sent } } is OutboxSendResult.NotDispatched -> { @@ -1654,12 +2148,73 @@ class ChatController internal constructor( } } + private suspend fun loadOutboxAttachmentsForSend( + outbox: ChatCommandOutbox, + item: ChatOutboxItem, + ): List { + if (item.attachments.isEmpty()) return emptyList() + return outbox.loadAttachments(item.id).map { loaded -> + OutgoingAttachment( + type = loaded.attachment.type, + mimeType = loaded.attachment.mimeType, + fileName = loaded.attachment.fileName, + base64 = Base64.getEncoder().encodeToString(loaded.bytes), + durationMs = loaded.attachment.durationMs, + ) + } + } + + /** + * Adopts run ownership for a flush-dispatched row in the visible session so streaming, the + * pending spinner, and reply reconciliation behave exactly like a direct send. The optimistic + * bubble replaces the queued row bubble until canonical history carries the turn. + */ + private fun adoptFlushedSend( + item: ChatOutboxItem, + attachments: List, + ackRunId: String, + ) { + if (normalizeRequestedSessionKey(item.sessionKey) != _sessionKey.value) return + val runId = item.id + if (locallyOwnedRun(runId) || locallyOwnedRun(ackRunId)) return + val optimistic = optimisticUserMessage(runId = runId, text = item.text, attachments = attachments) + optimisticMessagesByRunId[runId] = optimistic + unresolvedRepliesByRunId[runId] = optimistic + _messages.value = _messages.value + optimistic + armPendingRunTimeout(runId) + synchronized(pendingRuns) { + pendingRuns.add(runId) + _pendingRunCount.value = pendingRuns.size + } + // Chat events for this turn arrive under the acknowledged run id; mirroring the direct + // path's ownership transfer keeps the live run from looking foreign and timing out. + if (ackRunId != runId) transferRunOwnership(runId, ackRunId, optimistic) + } + private suspend fun attemptOutboxSend( + outbox: ChatCommandOutbox, item: ChatOutboxItem, gatewayId: String, + attachments: List, ): OutboxSendResult = try { val queuedSessionKey = normalizeRequestedSessionKey(item.sessionKey) + if (queuedSessionKey != item.sessionKey) { + // A row captured under the pre-hello "main" alias resolves exactly once, against the + // canonical main session active at first dispatch. Pinning it before the request means + // a later default-agent change can never redirect this input on a retry, so a pin + // that cannot be made durable must stop the dispatch while the row is still safe. + val pinned = + try { + outbox.pinSessionKey(item.id, queuedSessionKey) + true + } catch (err: CancellationException) { + throw err + } catch (_: Throwable) { + false + } + if (!pinned) return OutboxSendResult.NotDispatched("could not pin the delivery session") + } // Android only knows the active session's selected model. Unknown queued sessions fail // open, preserving the thinking level captured when they were enqueued. val thinking = @@ -1670,25 +2225,23 @@ class ChatController internal constructor( } else { item.thinkingLevel } + // The row id is the idempotency key, so gateway-side dedupe makes redelivery of an + // acked-but-crashed item harmless within the gateway's dedupe window. val params = - buildJsonObject { - // Rows enqueued under the pre-hello "main" alias must flush to the canonical main - // session the gateway announced, matching how the UI attributes those rows. - put("sessionKey", JsonPrimitive(queuedSessionKey)) - put("message", JsonPrimitive(item.text)) - put("thinking", JsonPrimitive(thinking)) - put("timeoutMs", JsonPrimitive(30_000)) - // The row id is the idempotency key, so gateway-side dedupe makes redelivery of an - // acked-but-crashed item harmless. - put("idempotencyKey", JsonPrimitive(item.id)) - } - val ack = parseChatSendAck(json, requestGatewayBound(gatewayId, "chat.send", params.toString())) + buildChatSendParams( + sessionKey = queuedSessionKey, + text = item.text, + thinking = thinking, + idempotencyKey = item.id, + attachments = attachments, + ) + val ack = parseChatSendAck(json, requestGatewayBound(gatewayId, "chat.send", params)) when (ack.normalizedStatus) { "ok", "started", "in_flight" -> if (ack.runId.isNullOrBlank()) { OutboxSendResult.DeliveryUnconfirmed(GatewayResponseState.Received) } else { - OutboxSendResult.Accepted + OutboxSendResult.Accepted(ack.runId) } "timeout", "error" -> OutboxSendResult.DeliveryUnconfirmed(GatewayResponseState.Received) else -> OutboxSendResult.DeliveryUnconfirmed(GatewayResponseState.Received) @@ -1969,9 +2522,30 @@ class ChatController internal constructor( terminalWithoutReplyRunIds.remove(runId) timedOutRunIds.add(runId) updateErrorText("Timed out waiting for a reply; try again or refresh.") + // The optimistic bubble is gone, so the journaled row must stay visible for review; + // history proof still retires it later if the turn did persist. + parkUnconfirmedDurableSend(runId) } } + /** Parks a still-accepted journaled row as delivery-unconfirmed once local ownership expires. */ + private suspend fun parkUnconfirmedDurableSend(runId: String) { + val outbox = commandOutbox ?: return + val row = + _outboxItems.value.firstOrNull { + it.status == ChatOutboxStatus.Accepted && + (it.id == runId || acknowledgedRunIdByRowId[it.id] == runId) + } ?: return + val persisted = updateOutboxStatusOrNull(outbox, row, ChatOutboxStatus.Failed, OUTBOX_DELIVERY_UNCONFIRMED_ERROR) + if (persisted == null) { + rearmOutboxRecovery() + _healthOk.value = false + } else { + acknowledgedRunIdByRowId.remove(row.id) + } + publishOutbox() + } + private fun clearPendingRun(runId: String) { pendingRunTimeoutJobs.remove(runId)?.cancel() unknownOutcomeRunIds.remove(runId) @@ -2140,6 +2714,10 @@ class ChatController internal constructor( unresolvedRunIds.forEach(unresolvedRepliesByRunId::remove) unresolvedRunIds.forEach(terminalWithoutReplyRunIds::remove) updateErrorText("Timed out confirming the sent message; refresh to check delivery.") + // Ownership expired without proof; keep the journaled copies visible for manual review. + for (unresolvedRunId in unresolvedRunIds) { + parkUnconfirmedDurableSend(unresolvedRunId) + } } } diff --git a/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatTranscriptCache.kt b/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatTranscriptCache.kt index 25987d417b2a..c6df2cd871e1 100644 --- a/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatTranscriptCache.kt +++ b/apps/android/app/src/main/java/ai/openclaw/app/chat/ChatTranscriptCache.kt @@ -158,8 +158,14 @@ internal interface ChatCacheDao { } @Database( - entities = [CachedSessionEntity::class, CachedMessageEntity::class, OutboxCommandEntity::class], - version = 3, + entities = [ + CachedSessionEntity::class, + CachedMessageEntity::class, + OutboxCommandEntity::class, + OutboxAttachmentEntity::class, + OutboxAttachmentChunkEntity::class, + ], + version = 4, exportSchema = false, ) internal abstract class ChatCacheDatabase : RoomDatabase() { @@ -186,13 +192,36 @@ internal abstract class ChatCacheDatabase : RoomDatabase() { } } + internal val MIGRATION_3_4 = + object : Migration(3, 4) { + override fun migrate(db: SupportSQLiteDatabase) { + db.execSQL("ALTER TABLE `outbox_commands` ADD COLUMN `gatedEpoch` INTEGER") + // Legacy queued command-shaped rows predate connection epochs; the sentinel makes + // them park for explicit retry instead of silently replaying on the next reconnect. + db.execSQL( + "UPDATE outbox_commands SET gatedEpoch = ? WHERE status = ? AND text LIKE '/%'", + arrayOf(OUTBOX_GATED_EPOCH_NEVER, ChatOutboxStatus.Queued.dbValue), + ) + db.execSQL( + "CREATE TABLE IF NOT EXISTS `outbox_attachments` (`id` TEXT NOT NULL, `commandId` TEXT NOT NULL, " + + "`position` INTEGER NOT NULL, `type` TEXT NOT NULL, `mimeType` TEXT NOT NULL, `fileName` TEXT NOT NULL, " + + "`durationMs` INTEGER, `byteLength` INTEGER NOT NULL, PRIMARY KEY(`id`))", + ) + db.execSQL("CREATE INDEX IF NOT EXISTS `index_outbox_attachments_commandId` ON `outbox_attachments` (`commandId`)") + db.execSQL( + "CREATE TABLE IF NOT EXISTS `outbox_attachment_chunks` (`attachmentId` TEXT NOT NULL, " + + "`chunkIndex` INTEGER NOT NULL, `bytes` BLOB NOT NULL, PRIMARY KEY(`attachmentId`, `chunkIndex`))", + ) + } + } + fun open( context: Context, name: String = CHAT_TRANSCRIPT_CACHE_DB_NAME, ): ChatCacheDatabase = Room .databaseBuilder(context, ChatCacheDatabase::class.java, name) - .addMigrations(MIGRATION_2_3) + .addMigrations(MIGRATION_2_3, MIGRATION_3_4) // v1 has only disposable transcripts. Starting with v2, the outbox is user data, so every // supported bump needs an explicit migration; destructive fallback remains for v1 only. .fallbackToDestructiveMigrationFrom(true, 1) diff --git a/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt b/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt index 2ea8649ccb92..eb4ba5ad0f4a 100644 --- a/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt +++ b/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatMessageViews.kt @@ -424,6 +424,7 @@ fun ChatOutboxBubble( when (item.status) { ChatOutboxStatus.Queued -> "Queued — sends when reconnected" ChatOutboxStatus.Sending -> "Sending…" + ChatOutboxStatus.Accepted -> "Sent — confirming delivery…" ChatOutboxStatus.Failed -> item.lastError ?.trim() @@ -435,7 +436,16 @@ fun ChatOutboxBubble( style = bubbleStyle("user").copy(borderColor = statusColor.copy(alpha = 0.6f)), roleLabel = "You", ) { - ChatMarkdown(text = item.text, textColor = mobileText) + if (item.text.isNotBlank()) { + ChatMarkdown(text = item.text, textColor = mobileText) + } + item.attachments.forEach { attachment -> + Text( + text = "📎 ${attachment.fileName}", + style = mobileCaption1, + color = mobileTextSecondary, + ) + } Row( verticalAlignment = Alignment.CenterVertically, horizontalArrangement = Arrangement.spacedBy(10.dp), @@ -449,7 +459,9 @@ fun ChatOutboxBubble( if (failed) { ChatOutboxAction(label = "Retry", color = mobileAccent, onClick = onRetry) } - if (item.status != ChatOutboxStatus.Sending) { + // Sending rows are mid-dispatch and accepted rows may already be delivered; both stay + // action-free until reconciliation resolves them, so a delete can never race a send. + if (item.status == ChatOutboxStatus.Queued || failed) { ChatOutboxAction(label = "Delete", color = mobileTextSecondary, onClick = onDelete) } } diff --git a/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt b/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt index da7b1a4863e8..601cc68b4a22 100644 --- a/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt +++ b/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatScreen.kt @@ -349,6 +349,7 @@ fun ChatScreen( items = outboxItems, sessionKey = sessionKey, mainSessionKey = mainSessionKey, + messages = messages, ), onRetryOutbox = viewModel::retryChatOutboxCommand, onDeleteOutbox = viewModel::deleteChatOutboxCommand, @@ -1143,15 +1144,13 @@ private fun ChatComposer( if (!thinkingSupported) thinkingSelectorExpanded = false } + // Offline sends queue durably too (text, images, and voice notes), so the gate is identical + // to the connected one; admission errors keep the draft when the durable queue refuses it. val sendEnabled = voiceNoteState !is VoiceNoteRecorderState.Recording && voiceNoteState !is VoiceNoteRecorderState.Preparing && pendingRunCount == 0 && - if (healthOk) { - value.trim().isNotEmpty() || attachments.isNotEmpty() - } else { - value.trim().isNotEmpty() && attachments.isEmpty() - } + (value.trim().isNotEmpty() || attachments.isNotEmpty()) Column(modifier = Modifier.fillMaxWidth().imePadding(), verticalArrangement = Arrangement.spacedBy(4.dp)) { if (attachments.isNotEmpty()) { @@ -1221,7 +1220,6 @@ private fun ChatComposer( ) } SendButton( - // Offline, only text sends are enabled: they queue durably (text-only v1). enabled = sendEnabled, onClick = onSend, ) diff --git a/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatTimeline.kt b/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatTimeline.kt index 2836ef7d4551..0ce45d53f53a 100644 --- a/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatTimeline.kt +++ b/apps/android/app/src/main/java/ai/openclaw/app/ui/chat/ChatTimeline.kt @@ -89,18 +89,25 @@ internal fun buildChatTimeline( /** * Outbox rows for the visible session. Rows enqueued under the "main" alias still belong to the - * canonical main session once the gateway hello rewrites the current key. + * canonical main session once the gateway hello rewrites the current key. Rows whose user turn + * is already visible as a message (optimistic while a live run owns it, or the canonical history + * copy right before the row retires) are hidden so one send never renders as two bubbles. */ internal fun outboxItemsForSession( items: List, sessionKey: String, mainSessionKey: String, + messages: List = emptyList(), ): List { val mainKey = mainSessionKey.trim().ifEmpty { "main" } val current = sessionKey.trim().let { if (it == "main") mainKey else it } + val visibleUserKeys = + messages + .mapNotNull { message -> message.idempotencyKey?.trim()?.takeIf { it.isNotEmpty() } } + .toSet() return items.filter { item -> val itemKey = item.sessionKey.let { if (it == "main") mainKey else it } - itemKey == current + itemKey == current && "${item.id}:user" !in visibleUserKeys } } diff --git a/apps/android/app/src/test/java/ai/openclaw/app/chat/ChatCacheDatabaseMigrationTest.kt b/apps/android/app/src/test/java/ai/openclaw/app/chat/ChatCacheDatabaseMigrationTest.kt index d77379f3b406..cfa5823626ac 100644 --- a/apps/android/app/src/test/java/ai/openclaw/app/chat/ChatCacheDatabaseMigrationTest.kt +++ b/apps/android/app/src/test/java/ai/openclaw/app/chat/ChatCacheDatabaseMigrationTest.kt @@ -4,6 +4,7 @@ import android.database.sqlite.SQLiteDatabase import kotlinx.coroutines.test.runTest import org.junit.Assert.assertEquals import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue import org.junit.Test import org.junit.runner.RunWith import org.robolectric.RobolectricTestRunner @@ -23,15 +24,17 @@ class ChatCacheDatabaseMigrationTest { val database = ChatCacheDatabase.open(context, databaseName) try { - // Opening through Room executes the production migration and validates the complete v3 - // schema, including columns, nullability, primary keys, defaults, and indices. - assertEquals(3, database.openHelper.writableDatabase.version) + // Opening through Room executes the production migration chain and validates the + // complete v4 schema, including columns, nullability, primary keys, and indices. + assertEquals(4, database.openHelper.writableDatabase.version) val outbox = RoomChatCommandOutbox(database) val rows = outbox.load("gateway-test").associateBy { it.id } val pristine = rows.getValue("pristine") assertEquals(ChatOutboxStatus.Queued, pristine.status) assertNull(pristine.lastError) + assertNull(pristine.gatedEpoch) + assertTrue(pristine.attachments.isEmpty()) for (id in listOf("legacy-queued-error", "interrupted-send")) { val migrated = rows.getValue(id) @@ -41,6 +44,11 @@ class ChatCacheDatabaseMigrationTest { val alreadyFailed = rows.getValue("already-failed") assertEquals(ChatOutboxStatus.Failed, alreadyFailed.status) assertEquals("original failure", alreadyFailed.lastError) + // Legacy queued command-shaped rows predate connection epochs; the sentinel keeps + // them from silently replaying on the next reconnect (they park for explicit retry). + val legacyCommand = rows.getValue("legacy-command") + assertEquals(ChatOutboxStatus.Queued, legacyCommand.status) + assertEquals(OUTBOX_GATED_EPOCH_NEVER, legacyCommand.gatedEpoch) assertEquals( "Cached session", database @@ -55,6 +63,49 @@ class ChatCacheDatabaseMigrationTest { } } + @Test + fun upgradedStoreSupportsAttachmentsAndSurvivesReopen() = + runTest { + val context = RuntimeEnvironment.getApplication() + val databaseName = "chat-cache-migration-${UUID.randomUUID()}.db" + val databaseFile = context.getDatabasePath(databaseName) + databaseFile.parentFile?.mkdirs() + createV2Fixture(databaseFile.path) + + // Spans multiple chunks to prove chunked reassembly is byte-exact after a real upgrade. + val bytes = ByteArray(OUTBOX_ATTACHMENT_CHUNK_BYTES + 77) { (it % 127).toByte() } + val queuedId: String + val first = ChatCacheDatabase.open(context, databaseName) + try { + val queued = + RoomChatCommandOutbox(first).enqueue( + gatewayId = "gateway-test", + sessionKey = "main", + text = "post-upgrade media", + thinkingLevel = "off", + nowMs = System.currentTimeMillis(), + attachments = + listOf( + OutboxAttachmentPayload(type = "image", mimeType = "image/jpeg", fileName = "a.jpg", durationMs = null, bytes = bytes), + ), + ) as ChatOutboxEnqueueResult.Queued + queuedId = queued.item.id + } finally { + first.close() + } + + // Process-restart analog: a fresh open must recover the exact bytes. + val reopened = ChatCacheDatabase.open(context, databaseName) + try { + val loaded = RoomChatCommandOutbox(reopened).loadAttachments(queuedId) + assertEquals(1, loaded.size) + assertTrue(bytes.contentEquals(loaded.single().bytes)) + } finally { + reopened.close() + context.deleteDatabase(databaseName) + } + } + private fun createV2Fixture(path: String) { SQLiteDatabase.openOrCreateDatabase(path, null).use { database -> val now = System.currentTimeMillis() @@ -105,6 +156,15 @@ class ChatCacheDatabaseMigrationTest { lastError = "original failure", createdAtMs = now + 3, ) + insertOutbox( + database, + id = "legacy-command", + status = "queued", + retryCount = 0, + lastError = null, + createdAtMs = now + 4, + text = "/clear", + ) database.version = 2 } } @@ -116,12 +176,13 @@ class ChatCacheDatabaseMigrationTest { retryCount: Int, lastError: String?, createdAtMs: Long, + text: String = id, ) { database.execSQL( "INSERT INTO outbox_commands " + "(id, gatewayId, sessionKey, text, thinkingLevel, createdAtMs, status, retryCount, lastError) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", - arrayOf(id, "gateway-test", "main", id, "off", createdAtMs, status, retryCount, lastError), + arrayOf(id, "gateway-test", "main", text, "off", createdAtMs, status, retryCount, lastError), ) } } diff --git a/apps/android/app/src/test/java/ai/openclaw/app/chat/ChatControllerOutboxTest.kt b/apps/android/app/src/test/java/ai/openclaw/app/chat/ChatControllerOutboxTest.kt index c421eacef1ac..132743fdb6f4 100644 --- a/apps/android/app/src/test/java/ai/openclaw/app/chat/ChatControllerOutboxTest.kt +++ b/apps/android/app/src/test/java/ai/openclaw/app/chat/ChatControllerOutboxTest.kt @@ -9,18 +9,28 @@ import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.launch +import kotlinx.coroutines.test.advanceTimeBy import kotlinx.coroutines.test.advanceUntilIdle import kotlinx.coroutines.test.runCurrent import kotlinx.coroutines.test.runTest import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonArray import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonPrimitive import org.junit.Assert.assertEquals import org.junit.Assert.assertFalse +import org.junit.Assert.assertNull import org.junit.Assert.assertTrue import org.junit.Test import java.util.UUID +private data class DeliveredSend( + val key: String, + val message: String, + val sessionKey: String, +) + @OptIn(ExperimentalCoroutinesApi::class) class ChatControllerOutboxTest { private val json = Json { ignoreUnknownKeys = true } @@ -36,13 +46,17 @@ class ChatControllerOutboxTest { private val capacity: Int = OUTBOX_MAX_QUEUED, ) : ChatCommandOutbox { val rows = LinkedHashMap() + val attachmentBytes = mutableMapOf>() val gatewayIds = mutableMapOf() val deletedSessions = mutableListOf() var recoveryGate: CompletableDeferred? = null var recoveryFailure: Throwable? = null var failedStatusUpdateFailure: Throwable? = null + var acceptedStatusUpdateFailure: Throwable? = null var queuedStatusUpdateFailure: Throwable? = null var sendingStatusUpdateFailure: Throwable? = null + var pinSessionKeyFailure: Throwable? = null + var claimGate: CompletableDeferred? = null var deleteFailure: Throwable? = null var deleteOnFailedStatus = false var loadGate: LoadGate? = null @@ -79,13 +93,22 @@ class ChatControllerOutboxTest { text: String, thinkingLevel: String, nowMs: Long, + attachments: List, + gatedEpoch: Long?, ): ChatOutboxEnqueueResult { if (gatewayIds.values.count { it == gatewayId } >= capacity) return ChatOutboxEnqueueResult.QueueFull + val commandBytes = attachments.sumOf { it.bytes.size.toLong() } + if (commandBytes > OUTBOX_MAX_COMMAND_ATTACHMENT_BYTES) return ChatOutboxEnqueueResult.AttachmentsTooLarge + val queuedBytes = attachmentBytes.values.sumOf { list -> list.sumOf { it.size.toLong() } } + if (commandBytes > 0 && queuedBytes + commandBytes > OUTBOX_MAX_GATEWAY_ATTACHMENT_BYTES) { + return ChatOutboxEnqueueResult.StorageFull + } val createdAt = maxOf(nowMs, nextCreatedAt) nextCreatedAt = createdAt + 1 + val id = UUID.randomUUID().toString() val item = ChatOutboxItem( - id = UUID.randomUUID().toString(), + id = id, sessionKey = sessionKey, text = text, thinkingLevel = thinkingLevel, @@ -93,12 +116,68 @@ class ChatControllerOutboxTest { status = ChatOutboxStatus.Queued, retryCount = 0, lastError = null, + gatedEpoch = gatedEpoch, + attachments = + attachments.mapIndexed { index, payload -> + ChatOutboxAttachment( + id = "$id-$index", + type = payload.type, + mimeType = payload.mimeType, + fileName = payload.fileName, + durationMs = payload.durationMs, + byteLength = payload.bytes.size.toLong(), + ) + }, ) rows[item.id] = item + attachmentBytes[item.id] = attachments.map { it.bytes } gatewayIds[item.id] = gatewayId return ChatOutboxEnqueueResult.Queued(item) } + override suspend fun loadAttachments(id: String): List { + val item = rows[id] ?: return emptyList() + val bytes = attachmentBytes[id].orEmpty() + return item.attachments.mapIndexed { index, attachment -> + LoadedOutboxAttachment(attachment = attachment, bytes = bytes[index]) + } + } + + override suspend fun claimForSending( + id: String, + retryCount: Int, + lastError: String?, + ): Int { + claimGate?.await() + sendingStatusUpdateFailure?.let { throw it } + val current = rows[id] ?: return 0 + if (current.status != ChatOutboxStatus.Queued) return 0 + rows[id] = current.copy(status = ChatOutboxStatus.Sending, retryCount = retryCount, lastError = lastError) + onStatusUpdated?.invoke(ChatOutboxStatus.Sending) + return 1 + } + + override suspend fun pinSessionKey( + id: String, + sessionKey: String, + ) { + pinSessionKeyFailure?.let { throw it } + val current = rows[id] ?: return + rows[id] = current.copy(sessionKey = sessionKey) + } + + override suspend fun confirmDelivered(ids: Set): Int { + var removed = 0 + for (id in ids) { + if (rows.remove(id) != null) { + attachmentBytes.remove(id) + gatewayIds.remove(id) + removed += 1 + } + } + return removed + } + override suspend fun updateStatus( id: String, status: ChatOutboxStatus, @@ -111,6 +190,7 @@ class ChatControllerOutboxTest { return 0 } if (status == ChatOutboxStatus.Failed) failedStatusUpdateFailure?.let { throw it } + if (status == ChatOutboxStatus.Accepted) acceptedStatusUpdateFailure?.let { throw it } if (status == ChatOutboxStatus.Queued) queuedStatusUpdateFailure?.let { throw it } if (status == ChatOutboxStatus.Sending) sendingStatusUpdateFailure?.let { throw it } val current = rows[id] ?: return 0 @@ -123,18 +203,41 @@ class ChatControllerOutboxTest { gatewayId: String, id: String, nowMs: Long, + gatedEpoch: Long?, ): Int { val current = rows[id] ?: return 0 if (gatewayIds[id] != gatewayId || current.status != ChatOutboxStatus.Failed) return 0 - val createdAt = maxOf(nowMs, nextCreatedAt) + var createdAt = maxOf(nowMs, nextCreatedAt) + rows[id] = + current.copy( + status = ChatOutboxStatus.Queued, + retryCount = 0, + lastError = null, + createdAtMs = createdAt, + gatedEpoch = gatedEpoch, + ) + // Mirror the Room store: queued same-session successors follow the retried row. + val successors = + rows.values + .filter { + it.id != id && + gatewayIds[it.id] == gatewayId && + it.sessionKey == current.sessionKey && + it.createdAtMs > current.createdAtMs && + it.status == ChatOutboxStatus.Queued + }.sortedBy { it.createdAtMs } + for (successor in successors) { + createdAt += 1 + rows[successor.id] = successor.copy(createdAtMs = createdAt) + } nextCreatedAt = createdAt + 1 - rows[id] = current.copy(status = ChatOutboxStatus.Queued, retryCount = 0, lastError = null, createdAtMs = createdAt) return 1 } override suspend fun delete(id: String) { deleteFailure?.let { throw it } rows.remove(id) + attachmentBytes.remove(id) gatewayIds.remove(id) } @@ -146,6 +249,7 @@ class ChatControllerOutboxTest { val ids = rows.values.filter { gatewayIds[it.id] == gatewayId && it.sessionKey == sessionKey }.map { it.id } ids.forEach { rows.remove(it) + attachmentBytes.remove(it) gatewayIds.remove(it) } } @@ -154,6 +258,7 @@ class ChatControllerOutboxTest { val ids = gatewayIds.filterValues { it == gatewayId }.keys.toList() ids.forEach { rows.remove(it) + attachmentBytes.remove(it) gatewayIds.remove(it) } } @@ -173,23 +278,35 @@ class ChatControllerOutboxTest { nowMs: Long, ) { for ((id, item) in rows) { - if (gatewayIds[id] == gatewayId && item.status == ChatOutboxStatus.Queued && item.createdAtMs <= nowMs - OUTBOX_EXPIRY_MS) { + if (gatewayIds[id] != gatewayId || item.createdAtMs > nowMs - OUTBOX_EXPIRY_MS) continue + if (item.status == ChatOutboxStatus.Queued) { rows[id] = item.copy(status = ChatOutboxStatus.Failed, lastError = OUTBOX_EXPIRED_ERROR) + } else if (item.status == ChatOutboxStatus.Accepted) { + rows[id] = item.copy(status = ChatOutboxStatus.Failed, lastError = OUTBOX_DELIVERY_UNCONFIRMED_ERROR) } } } } - /** Toggleable gateway seam: records chat.send idempotency keys and echoes them as run ids. */ + /** + * Toggleable gateway seam: records chat.send idempotency keys and echoes them as run ids. + * Sends that returned an acknowledgement are echoed into chat.history as `:user` rows + * plus an assistant reply, mirroring how the real gateway persists delivered turns; sends + * that threw after dispatch are not echoed (their persistence is genuinely unknown). + */ private inner class FakeGateway { var online = false var sendFailureBeforeDispatch: Throwable? = null var sendFailureAfterDispatch: Throwable? = null + var sendGate: CompletableDeferred? = null var sendResponse: (idempotencyKey: String) -> String = { key -> """{"runId":"$key","status":"started"}""" } val sentIdempotencyKeys = mutableListOf() val sentMessages = mutableListOf() val sentSessionKeys = mutableListOf() val sentThinkingLevels = mutableListOf() + val sentAttachmentFileNames = mutableListOf>() + var echoDeliveredSendsInHistory = true + private val deliveredSends = mutableListOf() var historyMessagesJson = "[]" var metadataModelsJson = "[]" @@ -203,14 +320,52 @@ class ChatControllerOutboxTest { sendFailureBeforeDispatch?.let { throw it } val params = json.parseToJsonElement(paramsJson.orEmpty()) as JsonObject val key = (params["idempotencyKey"] as JsonPrimitive).content + val message = (params["message"] as JsonPrimitive).content + val sessionKey = (params["sessionKey"] as JsonPrimitive).content sentIdempotencyKeys += key - sentMessages += (params["message"] as JsonPrimitive).content - sentSessionKeys += (params["sessionKey"] as JsonPrimitive).content + sentMessages += message + sentSessionKeys += sessionKey sentThinkingLevels += (params["thinking"] as JsonPrimitive).content + sentAttachmentFileNames += + (params["attachments"] as? JsonArray) + ?.mapNotNull { ((it as? JsonObject)?.get("fileName") as? JsonPrimitive)?.content } + .orEmpty() sendFailureAfterDispatch?.let { throw it } - sendResponse(key) + sendGate?.await() + val response = sendResponse(key) + // Terminal failures never persist a turn; every other returned ack means the gateway + // accepted the dispatch and the turn becomes visible in canonical history. + val status = + runCatching { (json.parseToJsonElement(response) as? JsonObject)?.get("status") as? JsonPrimitive } + .getOrNull() + ?.content + if (status != "timeout" && status != "error") { + deliveredSends += DeliveredSend(key = key, message = message, sessionKey = sessionKey) + } + response + } + "chat.history" -> { + val requestedKey = + runCatching { + ((json.parseToJsonElement(paramsJson.orEmpty()) as? JsonObject)?.get("sessionKey") as? JsonPrimitive)?.content + }.getOrNull() + val echoed = + if (echoDeliveredSendsInHistory) { + deliveredSends + .filter { requestedKey == null || it.sessionKey == requestedKey } + .flatMapIndexed { index, send -> + listOf( + """{"role":"user","content":"${send.message}","timestamp":${100 + index * 2},"idempotencyKey":"${send.key}:user"}""", + """{"role":"assistant","content":"reply","timestamp":${101 + index * 2},"idempotencyKey":"${send.key}:assistant"}""", + ) + } + } else { + emptyList() + } + val explicit = + (json.parseToJsonElement(historyMessagesJson) as JsonArray).map { it.toString() } + """{"sessionId":"session-1","messages":[${(explicit + echoed).joinToString(",")}]}""" } - "chat.history" -> """{"sessionId":"session-1","messages":$historyMessagesJson}""" "chat.metadata" -> """{"commands":[],"models":$metadataModelsJson}""" else -> "{}" } @@ -253,27 +408,6 @@ class ChatControllerOutboxTest { assertEquals(listOf("offline hello"), second.outboxItems.value.map { it.text }) } - @Test - fun offlineAttachmentSendsAreRejectedInsteadOfQueued() = - runTest { - val gateway = FakeGateway() - val outbox = FakeCommandOutbox() - val chat = controller(this, gateway, outbox) - chat.load("main") - advanceUntilIdle() - - val accepted = - chat.sendMessageAwaitAcceptance( - message = "with image", - thinkingLevel = "off", - attachments = listOf(OutgoingAttachment(type = "image", mimeType = "image/png", fileName = "a.png", base64 = "AAAA")), - ) - - assertFalse(accepted) - assertEquals("Gateway health not OK; cannot send", chat.errorText.value) - assertTrue(outbox.rows.isEmpty()) - } - @Test fun reconnectFlushesQueuedCommandsInOrderWithRowIdsAsIdempotencyKeys() = runTest { @@ -328,7 +462,7 @@ class ChatControllerOutboxTest { } @Test - fun failedAcceptedDeleteRearmsRecoveryBeforeYoungerRows() = + fun failedAcceptedPersistenceRearmsRecoveryBeforeYoungerRows() = runTest { val gateway = FakeGateway() val outbox = FakeCommandOutbox() @@ -346,7 +480,10 @@ class ChatControllerOutboxTest { attachments = emptyList(), ) - outbox.deleteFailure = IllegalStateException("storage unavailable") + // The acknowledged transition to accepted cannot be made durable; the flush must stop + // before younger rows instead of advancing past an ambiguous head still marked sending. + gateway.echoDeliveredSendsInHistory = false + outbox.acceptedStatusUpdateFailure = IllegalStateException("storage unavailable") gateway.online = true chat.handleGatewayEvent("health", null) advanceUntilIdle() @@ -366,15 +503,25 @@ class ChatControllerOutboxTest { .status, ) - outbox.deleteFailure = null + outbox.acceptedStatusUpdateFailure = null chat.handleGatewayEvent("health", null) - advanceUntilIdle() + // Bounded advance: enough for recovery, the flush, and its reconcile passes, but before + // the pending-run timeout would park the still-unproven younger send. + advanceTimeBy(5_000) + runCurrent() + // The re-armed recovery sweep parks the interrupted head for review and the younger row + // proceeds; the parked head no longer blocks the session. assertEquals(listOf("accepted", "younger"), gateway.sentMessages) - val recovered = chat.outboxItems.value.single() - assertEquals("accepted", recovered.text) - assertEquals(ChatOutboxStatus.Failed, recovered.status) - assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, recovered.lastError) + val parked = outbox.rows.values.first { it.text == "accepted" } + assertEquals(ChatOutboxStatus.Failed, parked.status) + assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, parked.lastError) + assertEquals( + ChatOutboxStatus.Accepted, + outbox.rows.values + .first { it.text == "younger" } + .status, + ) } @Test @@ -459,32 +606,6 @@ class ChatControllerOutboxTest { assertTrue(chat.outboxItems.value.isEmpty()) } - @Test - fun ackRemovesRowAndHistoryCopyIsTheOnlyBubble() = - runTest { - val gateway = FakeGateway() - val outbox = FakeCommandOutbox() - val chat = controller(this, gateway, outbox) - chat.load("main") - advanceUntilIdle() - - chat.sendMessageAwaitAcceptance(message = "queued text", thinkingLevel = "off", attachments = emptyList()) - val queuedRow = chat.outboxItems.value.single() - val queuedId = queuedRow.id - - // The post-flush history refresh returns the durable copy keyed by the row id. - gateway.historyMessagesJson = - """[{ "role": "user", "content": "queued text", "timestamp": 10, "idempotencyKey": "$queuedId" }]""" - gateway.online = true - chat.handleGatewayEvent("health", null) - advanceUntilIdle() - - assertTrue(chat.outboxItems.value.isEmpty()) - val userCopies = chat.messages.value.filter { message -> message.content.any { it.text == "queued text" } } - assertEquals(1, userCopies.size) - assertEquals(queuedId, userCopies.single().idempotencyKey) - } - @Test fun queuedRowsStayWithTheirGatewayAcrossSwitchAndFlushAfterSwitchBack() = runTest { @@ -1402,4 +1523,919 @@ class ChatControllerOutboxTest { assertEquals(listOf("agent:old:main"), outbox.deletedSessions) assertTrue(chat.outboxItems.value.isEmpty()) } + + @Test + fun offlineAttachmentSendsQueueDurablyWithByteRecovery() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + chat.load("main") + advanceUntilIdle() + + val imageBytes = byteArrayOf(1, 2, 3, 4) + val voiceBytes = byteArrayOf(9, 8, 7) + val accepted = + chat.sendMessageAwaitAcceptance( + message = "with media", + thinkingLevel = "off", + attachments = + listOf( + OutgoingAttachment( + type = "image", + mimeType = "image/jpeg", + fileName = "a.jpg", + base64 = + java.util.Base64 + .getEncoder() + .encodeToString(imageBytes), + ), + OutgoingAttachment( + type = "audio", + mimeType = "audio/mp4", + fileName = "note.m4a", + base64 = + java.util.Base64 + .getEncoder() + .encodeToString(voiceBytes), + durationMs = 1200L, + ), + ), + ) + + assertTrue(accepted) + val queued = chat.outboxItems.value.single() + assertEquals(ChatOutboxStatus.Queued, queued.status) + assertEquals(listOf("a.jpg", "note.m4a"), queued.attachments.map { it.fileName }) + assertEquals(1200L, queued.attachments[1].durationMs) + // Exact bytes survive the round trip into durable storage. + val loaded = outbox.loadAttachments(queued.id) + assertTrue(imageBytes.contentEquals(loaded[0].bytes)) + assertTrue(voiceBytes.contentEquals(loaded[1].bytes)) + + // Reconnect flushes the attachment payload with the captured metadata. + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + assertEquals(listOf(listOf("a.jpg", "note.m4a")), gateway.sentAttachmentFileNames) + assertTrue(chat.outboxItems.value.isEmpty()) + assertTrue(outbox.attachmentBytes.isEmpty()) + } + + @Test + fun historyProofRetiresRowAndTheCanonicalCopyIsTheOnlyBubble() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + chat.load("main") + advanceUntilIdle() + + chat.sendMessageAwaitAcceptance(message = "queued text", thinkingLevel = "off", attachments = emptyList()) + val queuedRow = chat.outboxItems.value.single() + val queuedId = queuedRow.id + + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + + assertTrue(chat.outboxItems.value.isEmpty()) + val userCopies = chat.messages.value.filter { message -> message.content.any { it.text == "queued text" } } + assertEquals(1, userCopies.size) + assertEquals("$queuedId:user", userCopies.single().idempotencyKey) + } + + @Test + fun acceptedRowSurvivesUntilCanonicalHistoryConfirmsIt() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + chat.load("main") + advanceUntilIdle() + + chat.sendMessageAwaitAcceptance(message = "await proof", thinkingLevel = "off", attachments = emptyList()) + val queuedId = + chat.outboxItems.value + .single() + .id + + // The gateway acks the send, but history lags: the durable row must not be deleted on + // the ACK alone, or a gateway crash before the transcript write would lose the message. + gateway.echoDeliveredSendsInHistory = false + gateway.online = true + chat.handleGatewayEvent("health", null) + runCurrent() + + assertEquals(listOf(queuedId), gateway.sentIdempotencyKeys) + assertEquals(ChatOutboxStatus.Accepted, outbox.rows.getValue(queuedId).status) + + // Canonical history catches up and retires the row. + gateway.echoDeliveredSendsInHistory = true + chat.refresh() + advanceUntilIdle() + assertFalse(outbox.rows.containsKey(queuedId)) + } + + @Test + fun healthySendsAreJournaledBeforeDispatchAndRetiredByHistoryProof() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.load("main") + advanceUntilIdle() + assertTrue(chat.healthOk.value) + + gateway.echoDeliveredSendsInHistory = false + val accepted = chat.sendMessageAwaitAcceptance(message = "healthy send", thinkingLevel = "off", attachments = emptyList()) + runCurrent() + + assertTrue(accepted) + // The dispatch used the durable row id as its idempotency key, and the row survives the + // started ACK: only canonical history proof may retire it. + val row = outbox.rows.values.single() + assertEquals(ChatOutboxStatus.Accepted, row.status) + assertEquals(listOf(row.id), gateway.sentIdempotencyKeys) + assertEquals(1, chat.messages.value.count { it.idempotencyKey == "${row.id}:user" }) + + gateway.echoDeliveredSendsInHistory = true + chat.refresh() + advanceUntilIdle() + assertTrue(outbox.rows.isEmpty()) + assertEquals(1, chat.messages.value.count { it.idempotencyKey == "${row.id}:user" }) + } + + @Test + fun processDeathDuringHealthyDispatchLeavesTheClaimForStartupRecovery() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val processJob = SupervisorJob() + val processScope = CoroutineScope(coroutineContext + processJob) + val first = controller(processScope, gateway, outbox) + gateway.online = true + first.load("main") + advanceUntilIdle() + + gateway.sendFailureAfterDispatch = CancellationException("process died mid-send") + runCatching { first.sendMessageAwaitAcceptance(message = "died in flight", thinkingLevel = "off", attachments = emptyList()) } + processJob.cancel() + + // The row keeps its 'sending' claim; the next process surfaces it as delivery-unconfirmed + // instead of silently replaying a possibly delivered dispatch. + assertEquals( + ChatOutboxStatus.Sending, + outbox.rows.values + .single() + .status, + ) + gateway.sendFailureAfterDispatch = null + gateway.echoDeliveredSendsInHistory = false + val restarted = controller(this, gateway, outbox) + restarted.handleGatewayEvent("health", null) + advanceUntilIdle() + val recovered = restarted.outboxItems.value.single() + assertEquals(ChatOutboxStatus.Failed, recovered.status) + assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, recovered.lastError) + assertEquals(1, gateway.sentMessages.size) + } + + @Test + fun restartOrphanedAcceptedRowIsRetiredByHistoryProofWithoutResending() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val processJob = SupervisorJob() + val processScope = CoroutineScope(coroutineContext + processJob) + val first = controller(processScope, gateway, outbox) + first.load("main") + advanceUntilIdle() + first.sendMessageAwaitAcceptance(message = "acked then killed", thinkingLevel = "off", attachments = emptyList()) + + // Flush accepts the row, then the process dies before history could confirm it. + gateway.echoDeliveredSendsInHistory = false + gateway.online = true + first.handleGatewayEvent("health", null) + runCurrent() + assertEquals( + ChatOutboxStatus.Accepted, + outbox.rows.values + .single() + .status, + ) + processJob.cancel() + + // The next process proves the turn against canonical history and retires the row + // without a second dispatch, even though the ACK was never locally processed further. + gateway.echoDeliveredSendsInHistory = true + val restarted = controller(this, gateway, outbox) + restarted.handleGatewayEvent("health", null) + advanceUntilIdle() + assertTrue(outbox.rows.isEmpty()) + assertEquals(1, gateway.sentIdempotencyKeys.size) + } + + @Test + fun restartOrphanedAcceptedRowWithoutHistoryProofParksForManualReview() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val processJob = SupervisorJob() + val processScope = CoroutineScope(coroutineContext + processJob) + val first = controller(processScope, gateway, outbox) + first.load("main") + advanceUntilIdle() + first.sendMessageAwaitAcceptance(message = "acked but lost", thinkingLevel = "off", attachments = emptyList()) + + gateway.echoDeliveredSendsInHistory = false + gateway.online = true + first.handleGatewayEvent("health", null) + runCurrent() + assertEquals( + ChatOutboxStatus.Accepted, + outbox.rows.values + .single() + .status, + ) + processJob.cancel() + + // The gateway lost the turn (crash between ACK and transcript write): an idle history + // without the row's key parks it for explicit review instead of auto-retrying. + val restarted = controller(this, gateway, outbox) + restarted.handleGatewayEvent("health", null) + advanceUntilIdle() + val parked = restarted.outboxItems.value.single() + assertEquals(ChatOutboxStatus.Failed, parked.status) + assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, parked.lastError) + assertEquals(1, gateway.sentIdempotencyKeys.size) + + // Explicit retry reuses the same idempotency key. + restarted.retryOutboxCommand(parked.id) + advanceUntilIdle() + assertEquals(listOf(parked.id, parked.id), gateway.sentIdempotencyKeys) + } + + @Test + fun preHelloMainRowsArePinnedAtFirstDispatchAndNeverRetarget() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + chat.load("main") + advanceUntilIdle() + chat.sendMessageAwaitAcceptance(message = "pinned input", thinkingLevel = "off", attachments = emptyList()) + assertEquals( + "main", + outbox.rows.values + .single() + .sessionKey, + ) + + // First dispatch resolves the alias against the hello-announced main session and pins it. + gateway.echoDeliveredSendsInHistory = false + gateway.sendFailureAfterDispatch = GatewayRequestOutcomeUnknown("ack lost") + gateway.online = true + chat.applyMainSessionKey("agent:work:main") + advanceUntilIdle() + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + val parked = outbox.rows.values.single() + assertEquals("agent:work:main", parked.sessionKey) + assertEquals(ChatOutboxStatus.Failed, parked.status) + + // A later default-agent change must not redirect the captured input on retry. + gateway.sendFailureAfterDispatch = null + chat.applyMainSessionKey("agent:other:main") + advanceUntilIdle() + chat.retryOutboxCommand(parked.id) + advanceUntilIdle() + assertEquals(listOf("agent:work:main", "agent:work:main"), gateway.sentSessionKeys) + } + + @Test + fun gatedCommandRowsParkAcrossReconnectAndSendOnlyOnExplicitRetry() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + var generation = 1L + val chat = + ChatController( + scope = this, + json = json, + requestGateway = gateway::request, + cacheScope = { ChatCacheScope(gatewayId = "gateway-test", connectionGeneration = generation) }, + commandOutbox = outbox, + ) + chat.load("main") + advanceUntilIdle() + + // A slash command captured offline is connection-gated to the epoch that captured it. + chat.sendMessageAwaitAcceptance(message = "/clear", thinkingLevel = "off", attachments = emptyList()) + val row = outbox.rows.values.single() + assertEquals(1L, row.gatedEpoch) + + // Reconnecting bumps the connection epoch, so the command parks instead of replaying. + generation = 2L + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + val parked = outbox.rows.values.single() + assertEquals(ChatOutboxStatus.Failed, parked.status) + assertEquals(OUTBOX_CONNECTION_CHANGED_ERROR, parked.lastError) + assertTrue(gateway.sentMessages.isEmpty()) + + // An explicit retry while connected re-arms the row for the live epoch and sends it. + chat.retryOutboxCommand(parked.id) + advanceUntilIdle() + assertEquals(listOf("/clear"), gateway.sentMessages) + assertTrue(chat.outboxItems.value.isEmpty()) + } + + @Test + fun directSlashSendParksWhenReconnectLandsBeforeDispatch() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + var generation = 1L + val chat = + ChatController( + scope = this, + json = json, + requestGateway = gateway::request, + cacheScope = { ChatCacheScope(gatewayId = "gateway-test", connectionGeneration = generation) }, + commandOutbox = outbox, + ) + gateway.online = true + chat.load("main") + advanceUntilIdle() + + // Hold the direct dispatch at its durable claim, then reconnect underneath it. The + // command was captured under epoch 1 and must not auto-send on the new connection. + outbox.claimGate = CompletableDeferred() + var accepted: Boolean? = null + val send = + launch { + accepted = chat.sendMessageAwaitAcceptance(message = "/clear", thinkingLevel = "off", attachments = emptyList()) + } + runCurrent() + generation = 2L + outbox.claimGate?.complete(Unit) + send.join() + advanceUntilIdle() + + assertEquals(true, accepted) + assertTrue(gateway.sentMessages.isEmpty()) + val parked = outbox.rows.values.single() + assertEquals(ChatOutboxStatus.Failed, parked.status) + assertEquals(OUTBOX_CONNECTION_CHANGED_ERROR, parked.lastError) + } + + @Test + fun acceptedRowAckedUnderDifferentRunIdStaysOwnedWhileTheRunIsLive() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.online = true + // The gateway acknowledges the send under a run id that differs from the row's + // idempotency key; local ownership transfers to that id while the row keeps its own. + gateway.sendResponse = { _ -> """{"runId":"gw-run-777","status":"started"}""" } + gateway.echoDeliveredSendsInHistory = false + chat.load("main") + advanceUntilIdle() + + val accepted = chat.sendMessageAwaitAcceptance(message = "slow turn", thinkingLevel = "off", attachments = emptyList()) + advanceTimeBy(1_000) + assertTrue(accepted) + assertEquals( + ChatOutboxStatus.Accepted, + outbox.rows.values + .single() + .status, + ) + + // A follow-up send must see the accepted head as live-owned: it dispatches directly, + // and the reconciliation sweep must not park the head while its run is in flight. + val followUp = chat.sendMessageAwaitAcceptance(message = "second", thinkingLevel = "off", attachments = emptyList()) + advanceTimeBy(10_000) + assertTrue(followUp) + assertEquals(listOf("slow turn", "second"), gateway.sentMessages) + assertTrue(outbox.rows.values.none { it.status == ChatOutboxStatus.Failed }) + } + + @Test + fun flushedSendAckedUnderDifferentRunIdResolvesWithTheLiveRun() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.sendResponse = { _ -> """{"runId":"gw-run-9","status":"started"}""" } + gateway.echoDeliveredSendsInHistory = false + chat.load("main") + advanceUntilIdle() + + // Captured offline, delivered by the reconnect flush under a divergent acked run id. + chat.sendMessageAwaitAcceptance(message = "queued turn", thinkingLevel = "off", attachments = emptyList()) + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceTimeBy(5_000) + assertEquals(listOf("queued turn"), gateway.sentMessages) + assertEquals( + ChatOutboxStatus.Accepted, + outbox.rows.values + .single() + .status, + ) + + // The run completes under the acknowledged id and its turn becomes visible in + // canonical history. The adopted send must resolve with the live run: without the + // ownership transfer the row-id pending run times out and surfaces a spurious error + // for a turn that was delivered. + gateway.echoDeliveredSendsInHistory = true + chat.handleGatewayEvent("chat", chatTerminalPayload("main", "gw-run-9", seq = 1, state = "final", assistantText = "done")) + advanceTimeBy(130_000) + assertEquals(0, chat.pendingRunCount.value) + assertTrue(outbox.rows.isEmpty()) + assertNull(chat.errorText.value) + } + + @Test + fun failedSessionPinKeepsTheRowQueuedInsteadOfDispatching() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + outbox.seed( + ChatOutboxItem( + id = "alias-row", + sessionKey = "main", + text = "captured pre-hello", + thinkingLevel = "off", + createdAtMs = System.currentTimeMillis(), + status = ChatOutboxStatus.Queued, + retryCount = 0, + lastError = null, + ), + ) + chat.load("main") + chat.applyMainSessionKey("agent:work:main") + advanceUntilIdle() + + // The durable pin is the only record of the alias resolution; if it cannot persist, + // dispatching anyway would let a retry after a default change target another session. + outbox.pinSessionKeyFailure = IllegalStateException("storage unavailable") + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + assertTrue(gateway.sentMessages.isEmpty()) + assertEquals(ChatOutboxStatus.Queued, outbox.rows.getValue("alias-row").status) + assertFalse(chat.healthOk.value) + + // Storage recovers; the next health transition pins and delivers exactly once. + outbox.pinSessionKeyFailure = null + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + assertEquals(listOf("agent:work:main"), gateway.sentSessionKeys) + assertTrue(outbox.rows.values.none { it.sessionKey == "main" }) + } + + @Test + fun reconcileParkWriteFailureFailsClosedThenParksAfterRecovery() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.echoDeliveredSendsInHistory = false + outbox.seed( + ChatOutboxItem( + id = "orphan-row", + sessionKey = "main", + text = "ambiguous send", + thinkingLevel = "off", + createdAtMs = System.currentTimeMillis(), + status = ChatOutboxStatus.Accepted, + retryCount = 0, + lastError = null, + ), + ) + chat.load("main") + advanceUntilIdle() + + // Two sightings without proof want to park the row, but the write fails: health drops + // instead of the reconciler claiming a change it never persisted. + outbox.failedStatusUpdateFailure = IllegalStateException("storage unavailable") + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + assertEquals(ChatOutboxStatus.Accepted, outbox.rows.getValue("orphan-row").status) + assertFalse(chat.healthOk.value) + + // Storage recovers; the next pass parks the orphan for manual review. + outbox.failedStatusUpdateFailure = null + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + val parked = outbox.rows.getValue("orphan-row") + assertEquals(ChatOutboxStatus.Failed, parked.status) + assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, parked.lastError) + } + + @Test + fun staleGatedParkFailureFailsClosedInsteadOfSpinning() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + outbox.seed( + ChatOutboxItem( + id = "stale-command", + sessionKey = "main", + text = "/clear", + thinkingLevel = "off", + createdAtMs = System.currentTimeMillis(), + status = ChatOutboxStatus.Queued, + retryCount = 0, + lastError = null, + gatedEpoch = 5L, + ), + ) + chat.load("main") + advanceUntilIdle() + + // The park write fails; the flush must drop health and stop instead of reloading the + // same stale row forever on a healthy connection. + outbox.failedStatusUpdateFailure = IllegalStateException("storage unavailable") + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + assertFalse(chat.healthOk.value) + assertEquals(ChatOutboxStatus.Queued, outbox.rows.getValue("stale-command").status) + assertTrue(gateway.sentMessages.isEmpty()) + + // Storage recovers; the next health transition parks the stale command for review. + outbox.failedStatusUpdateFailure = null + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + val parked = outbox.rows.getValue("stale-command") + assertEquals(ChatOutboxStatus.Failed, parked.status) + assertEquals(OUTBOX_CONNECTION_CHANGED_ERROR, parked.lastError) + assertTrue(gateway.sentMessages.isEmpty()) + } + + @Test + fun orphanedAcceptedHeadBlocksItsSessionUntilReconciliationParksIt() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val now = System.currentTimeMillis() + outbox.seed( + ChatOutboxItem( + id = "ambiguous-a", + sessionKey = "agent:a:main", + text = "unresolved head", + thinkingLevel = "off", + createdAtMs = now, + status = ChatOutboxStatus.Accepted, + retryCount = 0, + lastError = null, + ), + ) + outbox.seed( + ChatOutboxItem( + id = "queued-a", + sessionKey = "agent:a:main", + text = "blocked successor", + thinkingLevel = "off", + createdAtMs = now + 1, + status = ChatOutboxStatus.Queued, + retryCount = 0, + lastError = null, + ), + ) + outbox.seed( + ChatOutboxItem( + id = "queued-b", + sessionKey = "agent:b:main", + text = "independent session", + thinkingLevel = "off", + createdAtMs = now + 2, + status = ChatOutboxStatus.Queued, + retryCount = 0, + lastError = null, + ), + ) + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + + // The unproven accepted head held its session while the unrelated session flowed first; + // once reconciliation parked it for review, the released successor followed. + assertEquals(listOf("independent session", "blocked successor"), gateway.sentMessages) + assertEquals(ChatOutboxStatus.Failed, outbox.rows.getValue("ambiguous-a").status) + assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, outbox.rows.getValue("ambiguous-a").lastError) + assertFalse(outbox.rows.containsKey("queued-a")) + assertFalse(outbox.rows.containsKey("queued-b")) + } + + @Test + fun unconfirmedTimeoutParksTheAcceptedRowForReview() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.load("main") + advanceUntilIdle() + + // The gateway accepts the dispatch but its turn never reaches canonical history. + gateway.echoDeliveredSendsInHistory = false + chat.sendMessageAwaitAcceptance(message = "never confirmed", thinkingLevel = "off", attachments = emptyList()) + runCurrent() + assertEquals( + ChatOutboxStatus.Accepted, + outbox.rows.values + .single() + .status, + ) + + // Run ownership expires without proof; the row surfaces for manual review. + advanceUntilIdle() + val parked = outbox.rows.values.single() + assertEquals(ChatOutboxStatus.Failed, parked.status) + assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, parked.lastError) + } + + @Test + fun callerCancellationAfterTheClaimDoesNotStrandTheDirectSend() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.load("main") + advanceUntilIdle() + + // The UI scope dies (screen leaves composition) while the dispatch is suspended on the + // gateway response; the controller-owned dispatch must still settle the claimed row. + val gate = CompletableDeferred() + gateway.sendGate = gate + val callerJob = SupervisorJob() + val caller = CoroutineScope(coroutineContext + callerJob) + caller.launch { + chat.sendMessageAwaitAcceptance(message = "survives caller death", thinkingLevel = "off", attachments = emptyList()) + } + runCurrent() + assertEquals( + ChatOutboxStatus.Sending, + outbox.rows.values + .single() + .status, + ) + callerJob.cancel() + gate.complete(Unit) + advanceUntilIdle() + + // Delivered exactly once and retired by canonical history proof; nothing stranded. + assertEquals(listOf("survives caller death"), gateway.sentMessages) + assertTrue(outbox.rows.isEmpty()) + } + + @Test + fun directSendClaimFailureHandsDeliveryToTheFlushLane() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.load("main") + advanceUntilIdle() + + // The durable claim itself cannot be persisted; the admitted row must not be reported + // as sent-with-no-owner. The flush lane takes over and fails closed on the same error. + outbox.sendingStatusUpdateFailure = IllegalStateException("storage unavailable") + val accepted = chat.sendMessageAwaitAcceptance(message = "owned by flush", thinkingLevel = "off", attachments = emptyList()) + advanceUntilIdle() + assertTrue(accepted) + assertEquals( + ChatOutboxStatus.Queued, + outbox.rows.values + .single() + .status, + ) + assertTrue(gateway.sentMessages.isEmpty()) + assertFalse(chat.healthOk.value) + + // Storage recovers; the next health transition delivers the queued row exactly once. + outbox.sendingStatusUpdateFailure = null + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + assertEquals(listOf("owned by flush"), gateway.sentMessages) + assertTrue(outbox.rows.isEmpty()) + } + + @Test + fun directSendPersistenceFailureRearmsRecoveryInsteadOfStrandingSending() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.load("main") + advanceUntilIdle() + + // The acknowledged transition to accepted cannot be made durable mid-direct-send. + gateway.echoDeliveredSendsInHistory = false + outbox.acceptedStatusUpdateFailure = IllegalStateException("storage unavailable") + chat.sendMessageAwaitAcceptance(message = "stranded claim", thinkingLevel = "off", attachments = emptyList()) + runCurrent() + assertEquals( + ChatOutboxStatus.Sending, + outbox.rows.values + .single() + .status, + ) + assertFalse(chat.healthOk.value) + + // The re-armed recovery sweep parks the row on the next health transition, so the + // session is not blocked forever by a claim with no user action available. + outbox.acceptedStatusUpdateFailure = null + chat.handleGatewayEvent("health", null) + advanceTimeBy(5_000) + runCurrent() + val parked = outbox.rows.values.single() + assertEquals(ChatOutboxStatus.Failed, parked.status) + assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, parked.lastError) + } + + @Test + fun notEnqueuedDirectSendKeepsTheJournaledRowQueuedForReconnect() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.load("main") + advanceUntilIdle() + assertTrue(chat.healthOk.value) + + // The frame never enters the socket queue mid-direct-send; the durable copy must stay + // queued for reconnect instead of being deleted with only the volatile draft left. + gateway.sendFailureBeforeDispatch = GatewayRequestNotEnqueued("gateway send failed") + val accepted = chat.sendMessageAwaitAcceptance(message = "survives direct drop", thinkingLevel = "off", attachments = emptyList()) + assertTrue(accepted) + val row = outbox.rows.values.single() + assertEquals(ChatOutboxStatus.Queued, row.status) + assertFalse(chat.healthOk.value) + assertTrue(gateway.sentMessages.isEmpty()) + + gateway.sendFailureBeforeDispatch = null + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + assertEquals(listOf("survives direct drop"), gateway.sentMessages) + assertTrue(outbox.rows.isEmpty()) + } + + @Test + fun directDispatchWaitsForStartupRecoveryBeforeClaimingItsRow() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val recoveryGate = CompletableDeferred() + outbox.recoveryGate = recoveryGate + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.load("main") + runCurrent() + chat.setThinkingLevel("off") + + chat.sendMessage(message = "waits for recovery", thinkingLevel = "off", attachments = emptyList()) + runCurrent() + try { + // The row is journaled but must not be claimed 'sending' while the unscoped recovery + // sweep is pending, or the sweep would park this live dispatch as unconfirmed. + assertTrue(gateway.sentMessages.isEmpty()) + val row = outbox.rows.values.single() + assertEquals(ChatOutboxStatus.Queued, row.status) + } finally { + recoveryGate.complete(Unit) + } + advanceUntilIdle() + assertEquals(listOf("waits for recovery"), gateway.sentMessages) + assertTrue(outbox.rows.isEmpty()) + } + + @Test + fun ambiguousDirectSendKeepsTheComposerClearBecauseTheRowOwnsTheInput() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.load("main") + advanceUntilIdle() + + gateway.sendFailureAfterDispatch = IllegalStateException("transport wedged") + val accepted = chat.sendMessageAwaitAcceptance(message = "kept by the row", thinkingLevel = "off", attachments = emptyList()) + + // The dispatch outcome is unknown, so the journaled row parks for review and owns the + // input; a false return would restore a duplicate draft into the composer. + assertTrue(accepted) + val parked = outbox.rows.values.single() + assertEquals(ChatOutboxStatus.Failed, parked.status) + assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, parked.lastError) + assertEquals(1, gateway.sentMessages.size) + } + + @Test + fun historyProofOnABlockedHeadReleasesItsQueuedSuccessorInTheSameFlush() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val now = System.currentTimeMillis() + outbox.seed( + ChatOutboxItem( + id = "head", + sessionKey = "main", + text = "delivered before restart", + thinkingLevel = "off", + createdAtMs = now, + status = ChatOutboxStatus.Accepted, + retryCount = 0, + lastError = null, + ), + ) + outbox.seed( + ChatOutboxItem( + id = "tail", + sessionKey = "main", + text = "blocked successor", + thinkingLevel = "off", + createdAtMs = now + 1, + status = ChatOutboxStatus.Queued, + retryCount = 0, + lastError = null, + ), + ) + // Canonical history already carries the head's turn from the previous process. + gateway.historyMessagesJson = + """[{"role":"user","content":"delivered before restart","timestamp":5,"idempotencyKey":"head:user"},""" + + """{"role":"assistant","content":"r","timestamp":6,"idempotencyKey":"head:assistant"}]""" + val chat = controller(this, gateway, outbox) + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + + // Confirming the head must restart the drain so the released successor actually sends. + assertEquals(listOf("blocked successor"), gateway.sentMessages) + assertFalse(outbox.rows.containsKey("head")) + assertFalse(outbox.rows.containsKey("tail")) + } + + @Test + fun retryingAnUnconfirmedHeadWhileOfflineKeepsItAheadOfQueuedSuccessors() = + runTest { + val gateway = FakeGateway() + val outbox = FakeCommandOutbox() + val now = System.currentTimeMillis() + outbox.seed( + ChatOutboxItem( + id = "head", + sessionKey = "main", + text = "ambiguous head", + thinkingLevel = "off", + createdAtMs = now, + status = ChatOutboxStatus.Failed, + retryCount = 0, + lastError = OUTBOX_DELIVERY_UNCONFIRMED_ERROR, + ), + ) + outbox.seed( + ChatOutboxItem( + id = "tail", + sessionKey = "main", + text = "younger successor", + thinkingLevel = "off", + createdAtMs = now + 1, + status = ChatOutboxStatus.Queued, + retryCount = 0, + lastError = null, + ), + ) + val chat = controller(this, gateway, outbox) + chat.load("main") + advanceUntilIdle() + + // Retry while still offline: the head re-queues ahead of its still-queued successor, so + // the reconnect flush cannot deliver younger turns before the turn the user retried. + chat.retryOutboxCommand("head") + advanceUntilIdle() + gateway.online = true + chat.handleGatewayEvent("health", null) + advanceUntilIdle() + + assertEquals(listOf("ambiguous head", "younger successor"), gateway.sentMessages) + assertTrue(outbox.rows.isEmpty()) + } } diff --git a/apps/android/app/src/test/java/ai/openclaw/app/chat/RoomChatCommandOutboxTest.kt b/apps/android/app/src/test/java/ai/openclaw/app/chat/RoomChatCommandOutboxTest.kt index 9b878352f1b4..078a118a04cf 100644 --- a/apps/android/app/src/test/java/ai/openclaw/app/chat/RoomChatCommandOutboxTest.kt +++ b/apps/android/app/src/test/java/ai/openclaw/app/chat/RoomChatCommandOutboxTest.kt @@ -127,7 +127,7 @@ class RoomChatCommandOutboxTest { store.expireStale("gateway-a", nowMs = now) assertEquals(ChatOutboxStatus.Failed, store.load("gateway-a").single().status) - assertEquals(1, store.requeueForRetry(gatewayId = "gateway-a", id = stale.id, nowMs = now)) + assertEquals(1, store.requeueForRetry(gatewayId = "gateway-a", id = stale.id, nowMs = now, gatedEpoch = null)) store.expireStale("gateway-a", nowMs = now) val retried = store.load("gateway-a").single() @@ -143,7 +143,7 @@ class RoomChatCommandOutboxTest { val failed = store.enqueueQueued("gateway a failed", nowMs = 10, gatewayId = "gateway-a") store.updateStatus(failed.id, ChatOutboxStatus.Failed, retryCount = 1, lastError = "boom") - val changed = store.requeueForRetry(gatewayId = "gateway-b", id = failed.id, nowMs = 20) + val changed = store.requeueForRetry(gatewayId = "gateway-b", id = failed.id, nowMs = 20, gatedEpoch = null) assertEquals(0, changed) val untouched = store.load("gateway-a").single() @@ -157,11 +157,11 @@ class RoomChatCommandOutboxTest { runTest { val failed = store.enqueueQueued("retry once", nowMs = 10) store.updateStatus(failed.id, ChatOutboxStatus.Failed, retryCount = 1, lastError = "boom") - assertEquals(1, store.requeueForRetry(gatewayId = "gateway-a", id = failed.id, nowMs = 20)) + assertEquals(1, store.requeueForRetry(gatewayId = "gateway-a", id = failed.id, nowMs = 20, gatedEpoch = null)) store.updateStatus(failed.id, ChatOutboxStatus.Sending, retryCount = 0, lastError = null) val sendingCreatedAt = store.load("gateway-a").single().createdAtMs - val changed = store.requeueForRetry(gatewayId = "gateway-a", id = failed.id, nowMs = 30) + val changed = store.requeueForRetry(gatewayId = "gateway-a", id = failed.id, nowMs = 30, gatedEpoch = null) assertEquals(0, changed) val untouched = store.load("gateway-a").single() @@ -204,4 +204,220 @@ class RoomChatCommandOutboxTest { assertEquals(listOf("for other"), store.load("gateway-a").map { it.text }) } + + private fun payload( + bytes: ByteArray, + fileName: String = "a.jpg", + type: String = "image", + mimeType: String = "image/jpeg", + durationMs: Long? = null, + ): OutboxAttachmentPayload = OutboxAttachmentPayload(type = type, mimeType = mimeType, fileName = fileName, durationMs = durationMs, bytes = bytes) + + @Test + fun attachmentBytesRoundTripExactlyAcrossStoreReopen() = + runTest { + // Spans multiple chunks to prove chunked reassembly is byte-exact and ordered. + val big = ByteArray(OUTBOX_ATTACHMENT_CHUNK_BYTES + 1234) { (it % 251).toByte() } + val small = byteArrayOf(5, 4, 3) + val queued = + store.enqueue( + gatewayId = "gateway-a", + sessionKey = "main", + text = "with media", + thinkingLevel = "off", + nowMs = 10, + attachments = + listOf( + payload(big, fileName = "big.jpg"), + payload(small, fileName = "note.m4a", type = "audio", mimeType = "audio/mp4", durationMs = 900L), + ), + ) as ChatOutboxEnqueueResult.Queued + + val loadedItem = store.load("gateway-a").single() + assertEquals(listOf("big.jpg", "note.m4a"), loadedItem.attachments.map { it.fileName }) + assertEquals(listOf(big.size.toLong(), small.size.toLong()), loadedItem.attachments.map { it.byteLength }) + assertEquals(900L, loadedItem.attachments[1].durationMs) + + val loaded = store.loadAttachments(queued.item.id) + assertTrue(big.contentEquals(loaded[0].bytes)) + assertTrue(small.contentEquals(loaded[1].bytes)) + } + + @Test + fun perCommandAttachmentByteCapRefusesOversizedSends() = + runTest { + val oversized = ByteArray((OUTBOX_MAX_COMMAND_ATTACHMENT_BYTES + 1).toInt()) + val refused = + store.enqueue( + gatewayId = "gateway-a", + sessionKey = "main", + text = "too big", + thinkingLevel = "off", + nowMs = 10, + attachments = listOf(payload(oversized)), + ) + assertEquals(ChatOutboxEnqueueResult.AttachmentsTooLarge, refused) + assertTrue(store.load("gateway-a").isEmpty()) + } + + @Test + fun gatewayAttachmentByteBudgetRefusesWhenExhaustedAndRecoversAfterDelete() = + runTest { + val chunk = ByteArray(OUTBOX_MAX_COMMAND_ATTACHMENT_BYTES.toInt()) + val stored = mutableListOf() + var index = 0 + while (true) { + val result = + store.enqueue( + gatewayId = "gateway-a", + sessionKey = "main", + text = "bulk $index", + thinkingLevel = "off", + nowMs = index.toLong(), + attachments = listOf(payload(chunk)), + ) + if (result !is ChatOutboxEnqueueResult.Queued) { + assertEquals(ChatOutboxEnqueueResult.StorageFull, result) + break + } + stored += result.item.id + index += 1 + } + assertTrue(stored.isNotEmpty()) + + // Deleting a queued row releases its bytes, so admission recovers. + store.delete(stored.first()) + val retried = + store.enqueue( + gatewayId = "gateway-a", + sessionKey = "main", + text = "fits again", + thinkingLevel = "off", + nowMs = 999, + attachments = listOf(payload(chunk)), + ) + assertTrue(retried is ChatOutboxEnqueueResult.Queued) + } + + @Test + fun confirmDeliveredRetiresRowsAndTheirAttachmentBytesAtomically() = + runTest { + val bytes = byteArrayOf(1, 2, 3) + val queued = + store.enqueue( + gatewayId = "gateway-a", + sessionKey = "main", + text = "confirmed", + thinkingLevel = "off", + nowMs = 10, + attachments = listOf(payload(bytes)), + ) as ChatOutboxEnqueueResult.Queued + store.updateStatus(queued.item.id, ChatOutboxStatus.Accepted, retryCount = 0, lastError = null) + val keep = store.enqueueQueued("kept", nowMs = 20) + + assertEquals(1, store.confirmDelivered(setOf(queued.item.id, "missing-row"))) + + assertEquals(listOf(keep.id), store.load("gateway-a").map { it.id }) + assertTrue(store.loadAttachments(queued.item.id).isEmpty()) + } + + @Test + fun clearGatewayAndSessionDeleteAlsoDropAttachmentBytes() = + runTest { + val a = + store.enqueue( + gatewayId = "gateway-a", + sessionKey = "main", + text = "a", + thinkingLevel = "off", + nowMs = 10, + attachments = listOf(payload(byteArrayOf(1))), + ) as ChatOutboxEnqueueResult.Queued + val b = + store.enqueue( + gatewayId = "gateway-b", + sessionKey = "other", + text = "b", + thinkingLevel = "off", + nowMs = 20, + attachments = listOf(payload(byteArrayOf(2))), + ) as ChatOutboxEnqueueResult.Queued + + store.deleteForSession("gateway-b", "other") + store.clearGateway("gateway-a") + + assertTrue(store.load("gateway-a").isEmpty()) + assertTrue(store.load("gateway-b").isEmpty()) + assertTrue(store.loadAttachments(a.item.id).isEmpty()) + assertTrue(store.loadAttachments(b.item.id).isEmpty()) + } + + @Test + fun pinSessionKeyRewritesTheAliasExactlyOnce() = + runTest { + val queued = store.enqueueQueued("pinned", nowMs = 10) + store.pinSessionKey(queued.id, "agent:work:main") + assertEquals("agent:work:main", store.load("gateway-a").single().sessionKey) + } + + @Test + fun gatedEpochSurvivesPersistenceAndRetryRestamping() = + runTest { + val queued = + store.enqueue( + gatewayId = "gateway-a", + sessionKey = "main", + text = "/clear", + thinkingLevel = "off", + nowMs = 10, + gatedEpoch = 7L, + ) as ChatOutboxEnqueueResult.Queued + assertEquals(7L, store.load("gateway-a").single().gatedEpoch) + + store.updateStatus(queued.item.id, ChatOutboxStatus.Failed, retryCount = 0, lastError = OUTBOX_CONNECTION_CHANGED_ERROR) + assertEquals(1, store.requeueForRetry(gatewayId = "gateway-a", id = queued.item.id, nowMs = 20, gatedEpoch = 9L)) + assertEquals(9L, store.load("gateway-a").single().gatedEpoch) + } + + @Test + fun staleAcceptedRowsExpireToDeliveryUnconfirmed() = + runTest { + val now = 1_000_000_000L + val accepted = store.enqueueQueued("acked long ago", nowMs = now - OUTBOX_EXPIRY_MS - 1) + store.updateStatus(accepted.id, ChatOutboxStatus.Accepted, retryCount = 0, lastError = null) + + store.expireStale("gateway-a", nowMs = now) + + val row = store.load("gateway-a").single() + assertEquals(ChatOutboxStatus.Failed, row.status) + assertEquals(OUTBOX_DELIVERY_UNCONFIRMED_ERROR, row.lastError) + } + + @Test + fun claimForSendingIsAtomicAcrossCompetingDispatchers() = + runTest { + val queued = store.enqueueQueued("claim me", nowMs = 10) + + assertEquals(1, store.claimForSending(queued.id, 0, null)) + // The losing dispatcher gets 0 and must not send; the row is already claimed. + assertEquals(0, store.claimForSending(queued.id, 0, null)) + assertEquals(ChatOutboxStatus.Sending, store.load("gateway-a").single().status) + } + + @Test + fun requeueForRetryKeepsSameSessionQueuedSuccessorsBehindTheRetriedRow() = + runTest { + val head = store.enqueueQueued("head", nowMs = 10) + val tail = store.enqueueQueued("tail", nowMs = 20) + val other = store.enqueueQueued("other", nowMs = 30, sessionKey = "agent:other:main") + store.updateStatus(head.id, ChatOutboxStatus.Failed, retryCount = 0, lastError = OUTBOX_DELIVERY_UNCONFIRMED_ERROR) + + assertEquals(1, store.requeueForRetry(gatewayId = "gateway-a", id = head.id, nowMs = 1_000_000_000L, gatedEpoch = null)) + + val byId = store.load("gateway-a").associateBy { it.id } + // The retried head still precedes its session successor; unrelated sessions keep position. + assertTrue(byId.getValue(head.id).createdAtMs < byId.getValue(tail.id).createdAtMs) + assertEquals(ChatOutboxStatus.Queued, byId.getValue(tail.id).status) + assertEquals(30L, byId.getValue(other.id).createdAtMs) + } } diff --git a/apps/android/app/src/test/java/ai/openclaw/app/ui/chat/ChatTimelineTest.kt b/apps/android/app/src/test/java/ai/openclaw/app/ui/chat/ChatTimelineTest.kt index 3658aad98798..cba4668e54d6 100644 --- a/apps/android/app/src/test/java/ai/openclaw/app/ui/chat/ChatTimelineTest.kt +++ b/apps/android/app/src/test/java/ai/openclaw/app/ui/chat/ChatTimelineTest.kt @@ -2,6 +2,8 @@ package ai.openclaw.app.ui.chat import ai.openclaw.app.chat.ChatMessage import ai.openclaw.app.chat.ChatMessageContent +import ai.openclaw.app.chat.ChatOutboxItem +import ai.openclaw.app.chat.ChatOutboxStatus import ai.openclaw.app.chat.ChatPendingToolCall import org.junit.Assert.assertEquals import org.junit.Test @@ -88,6 +90,41 @@ class ChatTimelineTest { assertEquals(null, timeline.latestUserMessageId) } + @Test + fun outboxRowsHideOnceTheirUserTurnIsVisibleAsAMessage() { + val visible = + ChatOutboxItem( + id = "visible-row", + sessionKey = "main", + text = "still queued", + thinkingLevel = "off", + createdAtMs = 1, + status = ChatOutboxStatus.Queued, + retryCount = 0, + lastError = null, + ) + val consumed = + visible.copy( + id = "consumed-row", + status = ChatOutboxStatus.Accepted, + createdAtMs = 2, + ) + val optimisticCopy = + textMessage(id = "m1", role = "user", text = "sent already") + .copy(idempotencyKey = "consumed-row:user") + + val filtered = + outboxItemsForSession( + items = listOf(visible, consumed), + sessionKey = "main", + mainSessionKey = "agent:work:main", + messages = listOf(optimisticCopy), + ) + + // A row whose turn already renders as a message never shows a second bubble. + assertEquals(listOf("visible-row"), filtered.map { it.id }) + } + private fun textMessage( id: String, role: String, diff --git a/docs/platforms/android.md b/docs/platforms/android.md index 6e196ebdc80e..a3fb75594f1e 100644 --- a/docs/platforms/android.md +++ b/docs/platforms/android.md @@ -274,6 +274,7 @@ The Android Chat tab supports session selection (default `main`, plus other exis - History: `chat.history` (display-normalized — inline directive tags, plain-text tool-call XML payloads (``, ``, ``, ``, and truncated variants), and leaked ASCII/full-width model control tokens are stripped; silent-token assistant rows such as exact `NO_REPLY` / `no_reply` are omitted; oversized rows can be replaced with placeholders) - Send: `chat.send` +- Durable sending: every send (text, picked images, and voice notes) is journaled to a per-gateway on-device outbox before any network attempt, so app termination cannot lose submitted input. Sends queued while offline deliver in order on reconnect with stable idempotency keys, and a send is retired only after the turn is visible in canonical `chat.history` — an acknowledgement alone is not treated as proof of delivery. Ambiguous outcomes (lost acknowledgement, app killed mid-send, gateway restart before the transcript write) surface as visible rows with explicit **Retry**/**Delete** instead of auto-resending. Slash commands never auto-replay across a reconnect; they park for explicit retry. The queue is bounded (50 messages and 48 MB of attachment bytes per gateway) and unsent rows expire after 48 hours. Composer drafts that were never submitted are not process-durable. - Push updates (best-effort): `chat.subscribe` -> `event:"chat"` - Listen: long-press an assistant message and choose **Listen** to hear it; audio renders via gateway `tts.speak` with the configured TTS provider chain, and on-device system TTS is used when the gateway cannot render audio. Playback stops on session switch, new chat, app backgrounding, or chat close.