Fixed tool calling history
This commit is contained in:
1 parent
dd43033888
commit
524ed8ce53
11 files changed
+471
-55
No files matched your search
@@ -16,9 +16,9 @@ import iris.protocol.IrisJson
|
||||
* the controller) so every reconciled frame survives a process death.
|
||||
* - [metaGet] / [metaPut] hold small UI state (last-viewed lane).
|
||||
*
|
||||
* Only [MessageItem]s are persisted — tool cards, live streaming state and
|
||||
* [MessageItem]s and [ToolItem]s are persisted; live streaming state and
|
||||
* local system notices are ephemeral. Rows are JSON payloads keyed by
|
||||
* (lane, id), so the schema does not drift with [MessageItem] fields.
|
||||
* (lane, id), so the schema does not drift with the model fields.
|
||||
*
|
||||
* All access is synchronized: the debounced save collectors run on the
|
||||
* controller scope while [dispose] may flush from the UI thread.
|
||||
@@ -32,27 +32,70 @@ class ChatDb(
|
||||
|
||||
// ── messages ──────────────────────────────────────────────────────────
|
||||
|
||||
/** All persisted lanes (lane key -> messages ordered by ts). */
|
||||
fun loadLanes(): Map<String, List<MessageItem>> =
|
||||
/** All persisted lanes (lane key -> items: messages ordered by ts with
|
||||
* tool cards interleaved at their anchored position). */
|
||||
fun loadLanes(): Map<String, List<ChatItem>> =
|
||||
synchronized(lock) {
|
||||
val lanes = linkedMapOf<String, MutableList<MessageItem>>()
|
||||
val messages = linkedMapOf<String, MutableList<MessageItem>>()
|
||||
for (row in db.cacheQueries.allMessages().executeAsList()) {
|
||||
val item = decodeMessage(row.payload) ?: continue
|
||||
lanes.getOrPut(row.lane) { mutableListOf() }.add(item)
|
||||
messages.getOrPut(row.lane) { mutableListOf() }.add(item)
|
||||
}
|
||||
val tools = linkedMapOf<String, MutableList<ToolItem>>()
|
||||
for (row in db.cacheQueries.allTools().executeAsList()) {
|
||||
val item = decodeTool(row.payload) ?: continue
|
||||
tools.getOrPut(row.lane) { mutableListOf() }.add(item)
|
||||
}
|
||||
(messages.keys + tools.keys).distinct().associateWith { lane ->
|
||||
val items: MutableList<ChatItem> = messages[lane].orEmpty().toMutableList()
|
||||
// Insert each tool card after its anchor message (the message
|
||||
// it followed live). Cards sharing an anchor keep their seq
|
||||
// order; a card whose anchor is gone (deleted message) falls
|
||||
// to the end of the lane. The lastPos cache assumes anchors
|
||||
// are monotonically non-decreasing per lane (true for the
|
||||
// onToolStart anchor rule: last non-streaming message) — a
|
||||
// tool anchored to an EARLIER message processed after a
|
||||
// later-anchored one would be misplaced.
|
||||
val lastPos = mutableMapOf<String, Int>()
|
||||
for (tool in tools[lane].orEmpty()) {
|
||||
val anchor = tool.anchorId
|
||||
val pos =
|
||||
if (anchor == null) {
|
||||
-1
|
||||
} else {
|
||||
lastPos.getOrPut(anchor) { items.indexOfFirst { it.id == anchor } }
|
||||
}
|
||||
if (pos >= 0) {
|
||||
items.add(pos + 1, tool)
|
||||
if (anchor != null) lastPos[anchor] = pos + 1
|
||||
} else {
|
||||
items.add(tool)
|
||||
}
|
||||
}
|
||||
items
|
||||
}
|
||||
lanes.mapValues { it.value.toList() }
|
||||
}
|
||||
|
||||
/** Replace the whole message cache with [lanes] (atomic snapshot). Tool
|
||||
* cards and local system notices are skipped (ephemeral). */
|
||||
/** Replace the whole message + tool cache with [lanes] (atomic snapshot).
|
||||
* Local system notices are skipped (ephemeral). */
|
||||
fun saveLanes(lanes: Map<String, List<ChatItem>>) {
|
||||
synchronized(lock) {
|
||||
db.transaction {
|
||||
db.cacheQueries.clearMessages()
|
||||
db.cacheQueries.clearTools()
|
||||
for ((lane, items) in lanes) {
|
||||
var toolSeq = 0
|
||||
for (item in items) {
|
||||
if (item is MessageItem && !item.isSystem) {
|
||||
db.cacheQueries.upsertMessage(lane, item.id, item.ts, json.encodeToString(item))
|
||||
when (item) {
|
||||
is MessageItem -> {
|
||||
if (!item.isSystem) {
|
||||
db.cacheQueries.upsertMessage(lane, item.id, item.ts, json.encodeToString(item))
|
||||
}
|
||||
}
|
||||
|
||||
is ToolItem -> {
|
||||
db.cacheQueries.upsertTool(lane, item.id, toolSeq++.toLong(), json.encodeToString(item))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -109,6 +152,7 @@ class ChatDb(
|
||||
synchronized(lock) {
|
||||
db.transaction {
|
||||
db.cacheQueries.clearMessages()
|
||||
db.cacheQueries.clearTools()
|
||||
db.cacheQueries.clearChannels()
|
||||
db.cacheQueries.clearMeta()
|
||||
}
|
||||
@@ -122,6 +166,13 @@ class ChatDb(
|
||||
null
|
||||
}
|
||||
|
||||
private fun decodeTool(payload: String): ToolItem? =
|
||||
try {
|
||||
json.decodeFromString<ToolItem>(payload).sanitizeForRestore()
|
||||
} catch (_: Exception) {
|
||||
null
|
||||
}
|
||||
|
||||
/** A restored message is never mid-flight: a streaming bubble is
|
||||
* finalized (the sync replay finalizes it for real on reconnect) and a
|
||||
* pending send becomes failed (tap to retry) — the gateway never
|
||||
@@ -132,4 +183,9 @@ class ChatDb(
|
||||
pending = false,
|
||||
status = if (status == MsgStatus.Pending) MsgStatus.Failed else status,
|
||||
)
|
||||
|
||||
/** A restored tool card is never mid-flight: an open card (the process
|
||||
* died before tool.end) is closed as interrupted, mirroring
|
||||
* [ChatStore.finalizeInterrupted]. */
|
||||
private fun ToolItem.sanitizeForRestore(): ToolItem = if (done) this else copy(done = true, ok = false)
|
||||
}
|
||||
@@ -90,7 +90,13 @@ data class MediaItem(
|
||||
val localPath: String? = null,
|
||||
)
|
||||
|
||||
/** A structured tool-activity card (spinner until [done]). */
|
||||
/** A structured tool-activity card (spinner until [done]).
|
||||
* [anchorId] is the id of the message this card follows in the lane (the
|
||||
* last non-streaming message when the tool started) — persisted with the
|
||||
* card so a restart restores it in its correct position (user message →
|
||||
* tool card → answer) instead of dropping it or appending it at the end.
|
||||
* [Serializable]: persisted as a JSON payload in the local cache (ChatDb). */
|
||||
@Serializable
|
||||
data class ToolItem(
|
||||
override val id: String,
|
||||
val index: Int,
|
||||
@@ -102,6 +108,7 @@ data class ToolItem(
|
||||
val ok: Boolean = true,
|
||||
val duration: Double? = null,
|
||||
val outputPreview: String? = null,
|
||||
val anchorId: String? = null,
|
||||
) : ChatItem
|
||||
|
||||
class ChatStore {
|
||||
@@ -114,7 +121,6 @@ class ChatStore {
|
||||
val currentLane: StateFlow<String> = _currentLane.asStateFlow()
|
||||
|
||||
private var localSeq = 0
|
||||
private var toolSeq = 0
|
||||
|
||||
/** When false, `message.start`/`message.update` frames are ignored and each
|
||||
* reply materializes as a single final message on `message.stop`
|
||||
@@ -169,8 +175,11 @@ class ChatStore {
|
||||
lane: String,
|
||||
media: List<MediaItem> = emptyList(),
|
||||
): String {
|
||||
localSeq++
|
||||
val id = "local_$localSeq"
|
||||
// Process-unique id: the in-memory seq resets on every ChatStore
|
||||
// creation, and failed sends are persisted — a restart would
|
||||
// otherwise re-mint local_1 and collide with the restored bubble
|
||||
// (duplicate list key).
|
||||
val id = randomId("local")
|
||||
updateLane(lane) {
|
||||
it + MessageItem(id = id, role = ROLE_USER, text = text, ts = 0, pending = true, status = MsgStatus.Pending, media = media)
|
||||
}
|
||||
@@ -418,9 +427,17 @@ class ChatStore {
|
||||
frame: Frame,
|
||||
) {
|
||||
val p = frame.payloadAs<ToolStartPayload>() ?: return
|
||||
toolSeq++
|
||||
val id = "tool_$toolSeq"
|
||||
// Process-unique id: the in-memory seq resets on every ChatStore
|
||||
// creation, and cards are persisted — a restart would otherwise
|
||||
// re-mint tool_1 and collide with the restored card (duplicate list
|
||||
// key, upsert overwrite). onToolProgress/onToolEnd match by index,
|
||||
// so non-sequential ids are safe.
|
||||
val id = randomId("tool")
|
||||
updateLane(lane) { list ->
|
||||
// Anchor the card to the message it follows: the last non-streaming
|
||||
// message (a live streaming bubble is the answer that arrives AFTER
|
||||
// the tool, so it is skipped). Restored with the card on restart.
|
||||
val anchorId = list.lastOrNull { it is MessageItem && !it.streaming }?.id
|
||||
list +
|
||||
ToolItem(
|
||||
id = id,
|
||||
@@ -428,6 +445,7 @@ class ChatStore {
|
||||
name = p.name,
|
||||
preview = p.preview,
|
||||
args = p.args,
|
||||
anchorId = anchorId,
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -682,9 +700,11 @@ class ChatStore {
|
||||
* Load a history page into [lane] (oldest → newest). The history is the
|
||||
* authoritative full list of final messages for the lane; it replaces the
|
||||
* lane's final messages and preserves non-final items (tool cards, live
|
||||
* streaming bubbles, commentary) that are not part of the history. Used to
|
||||
* restore the view on first open of a chat / after a process death, where
|
||||
* the in-memory store is empty and the `sync` delta does not cover older
|
||||
* streaming bubbles, commentary) that are not part of the history. Tool
|
||||
* cards sort at their anchor message's position, so a history refresh
|
||||
* keeps them between the user message and the answer. Used to restore the
|
||||
* view on first open of a chat / after a process death, where the
|
||||
* in-memory store is empty and the `sync` delta does not cover older
|
||||
* messages.
|
||||
*/
|
||||
fun loadHistory(
|
||||
@@ -697,8 +717,28 @@ class ChatStore {
|
||||
list.filter { item ->
|
||||
item !is MessageItem || item.id !in historyIds
|
||||
}
|
||||
// ts of every item in the current lane: ts-less items (commentary,
|
||||
// tool cards) inherit the ts of the item before them, so a tool
|
||||
// card anchored to a commentary still sorts at the right place.
|
||||
val tsOf = mutableMapOf<String, Long>()
|
||||
var lastTs = 0L
|
||||
for (item in list) {
|
||||
val t = (item as? MessageItem)?.ts?.takeIf { it > 0 } ?: lastTs
|
||||
if (t > 0) lastTs = t
|
||||
tsOf[item.id] = t
|
||||
}
|
||||
// History is authoritative for the ts of its messages.
|
||||
for (m in messages) {
|
||||
if (m.ts > 0) tsOf[m.id] = m.ts
|
||||
}
|
||||
(messages + preserved).sortedBy { item ->
|
||||
(item as? MessageItem)?.ts?.takeIf { it > 0 } ?: Long.MAX_VALUE
|
||||
when (item) {
|
||||
is MessageItem -> item.ts.takeIf { it > 0 } ?: Long.MAX_VALUE
|
||||
|
||||
// A resolved ts of 0 means the anchor itself is ts-less
|
||||
// (lane start) — sort with it (end) instead of to the top.
|
||||
is ToolItem -> tsOf[item.anchorId]?.takeIf { it > 0 } ?: Long.MAX_VALUE
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -710,7 +750,7 @@ class ChatStore {
|
||||
* delta and the `history` refresh reconcile the cache on connect.
|
||||
* No-op when [lanes] is empty (first launch).
|
||||
*/
|
||||
fun loadFromCache(lanes: Map<String, List<MessageItem>>) {
|
||||
fun loadFromCache(lanes: Map<String, List<ChatItem>>) {
|
||||
if (lanes.isEmpty()) return
|
||||
_lanes.value = lanes
|
||||
}
|
||||
|
||||
@@ -465,13 +465,24 @@ class IrisController(
|
||||
client.events.collect { frame ->
|
||||
try {
|
||||
when (frame.type) {
|
||||
TYPE_TOOL_START,
|
||||
TYPE_TOOL_PROGRESS,
|
||||
TYPE_TOOL_END,
|
||||
-> {
|
||||
// Tool cards are restored from the local cache
|
||||
// (anchored to their message); they are not part
|
||||
// of `history`. A sync replay (frame carries a
|
||||
// cursor) would create duplicate cards appended
|
||||
// AFTER the lane's restored messages — drop
|
||||
// them. Live tool frames (no cursor) flow through
|
||||
// as usual.
|
||||
if (frame.cursor == null) chat.onFrame(frame)
|
||||
}
|
||||
|
||||
TYPE_MESSAGE,
|
||||
TYPE_MESSAGE_START,
|
||||
TYPE_MESSAGE_UPDATE,
|
||||
TYPE_MESSAGE_STOP,
|
||||
TYPE_TOOL_START,
|
||||
TYPE_TOOL_PROGRESS,
|
||||
TYPE_TOOL_END,
|
||||
TYPE_COMMENTARY,
|
||||
TYPE_MEDIA_OFFER,
|
||||
-> {
|
||||
|
||||
@@ -11,6 +11,14 @@ CREATE TABLE message (
|
||||
PRIMARY KEY (lane, id)
|
||||
);
|
||||
|
||||
CREATE TABLE tool (
|
||||
lane TEXT NOT NULL, -- lane key: chatId or chatId::threadId
|
||||
id TEXT NOT NULL, -- local tool card id (tool_N)
|
||||
seq INTEGER NOT NULL, -- card order within the lane (lane position)
|
||||
payload TEXT NOT NULL, -- serialized ToolItem (carries its anchor_id)
|
||||
PRIMARY KEY (lane, id)
|
||||
);
|
||||
|
||||
CREATE TABLE channel (
|
||||
chat_id TEXT NOT NULL PRIMARY KEY,
|
||||
payload TEXT NOT NULL -- serialized ChannelInfo
|
||||
@@ -33,6 +41,18 @@ VALUES (?, ?, ?, ?);
|
||||
clearMessages:
|
||||
DELETE FROM message;
|
||||
|
||||
allTools:
|
||||
SELECT lane, seq, payload
|
||||
FROM tool
|
||||
ORDER BY lane, seq;
|
||||
|
||||
upsertTool:
|
||||
INSERT OR REPLACE INTO tool (lane, id, seq, payload)
|
||||
VALUES (?, ?, ?, ?);
|
||||
|
||||
clearTools:
|
||||
DELETE FROM tool;
|
||||
|
||||
allChannels:
|
||||
SELECT chat_id, payload
|
||||
FROM channel;
|
||||
|
||||
@@ -0,0 +1,9 @@
|
||||
-- v1 -> v2: persist tool cards (M-cache: tool cards survive a restart,
|
||||
-- anchored to the message they follow — see ChatDb.loadLanes).
|
||||
CREATE TABLE tool (
|
||||
lane TEXT NOT NULL,
|
||||
id TEXT NOT NULL,
|
||||
seq INTEGER NOT NULL,
|
||||
payload TEXT NOT NULL,
|
||||
PRIMARY KEY (lane, id)
|
||||
);
|
||||
Reference in new issue
Block a user