| 1 | import type { Action, Item, State } from "./useController"; |
| 2 | import { sessionIdentityStableKey } from "./sessionIdentity"; |
| 3 | import { beginLocalSubmission, canonicalUserConfirmations, orderedLocalSubmissions, settleLocalSubmissions, updateLocalSubmission, type LocalSubmission } from "./localSubmissionState"; |
| 4 | import { recordFrontendDiagnostic } from "./frontendDiagnosticBridge"; |
| 5 | |
| 6 | export function submissionBindingCurrent(current: State | undefined, expected: State): boolean { |
| 7 | return current?.sessionGen === expected.sessionGen && sessionIdentityStableKey(current?.meta) === sessionIdentityStableKey(expected.meta); |
| 8 | } |
| 9 | |
| 10 | export function resetTurnTiming(now = Date.now()): Pick<State, "turnStartAt" | "turnDoneAt" | "turnWaitAccumMs" | "promptWaitStartedAt" | "turnTokens" | "turnTotalTokens" | "turnUsage" | "turnOutputTokens" | "turnOutputChars" | "turnOutputCharsAtUsage" | "turnOutputEstimated" | "turnModelActiveAt" | "turnModelActiveMs" | "turnRateSample" | "turnCost" | "turnRateBand" | "turnArgChars" | "pendingRequestModelMs"> { |
| 11 | return { |
| 12 | turnStartAt: now, |
| 13 | turnDoneAt: 0, |
| 14 | turnWaitAccumMs: 0, |
| 15 | promptWaitStartedAt: undefined, |
| 16 | turnTokens: 0, |
| 17 | turnTotalTokens: 0, |
| 18 | turnUsage: undefined, |
| 19 | turnOutputTokens: 0, |
| 20 | turnOutputChars: 0, |
| 21 | turnOutputCharsAtUsage: 0, |
| 22 | turnOutputEstimated: false, |
| 23 | turnModelActiveAt: undefined, |
| 24 | turnModelActiveMs: 0, pendingRequestModelMs: undefined, |
| 25 | turnRateSample: undefined, |
| 26 | turnCost: 0, |
| 27 | turnRateBand: undefined, |
| 28 | turnArgChars: 0, |
| 29 | }; |
| 30 | } |
| 31 | |
| 32 | export function confirmPendingUser(s: State, submissionId: string | undefined): State { |
| 33 | if (!submissionId) return s; |
| 34 | const next = s.pendingSubmissionId === submissionId ? { ...s, pendingUser: undefined, pendingSubmissionId: undefined } : s; |
| 35 | if (s.localSubmissions[submissionId]?.status === "failed") return next; |
| 36 | return updateLocalSubmission(next, submissionId, { |
| 37 | status: "accepted", |
| 38 | }); |
| 39 | } |
| 40 | |
| 41 | |
| 42 | function retainResidentTurnId(update: Item, resident: Item): Item { |
| 43 | if (update.turnId || !resident.turnId) return update; |
| 44 | return { ...update, turnId: resident.turnId }; |
| 45 | } |
| 46 | |
| 47 | function submissionForInstalledUser(s: State, item: Item): LocalSubmission | undefined { |
| 48 | if (item.kind !== "user") return undefined; |
| 49 | const direct = item.submissionId ? s.localSubmissions[item.submissionId] : undefined; |
| 50 | if (direct && (!direct.messageId || !item.messageId || direct.messageId === item.messageId)) return direct; |
| 51 | if (!item.messageId) return undefined; |
| 52 | return orderedLocalSubmissions(s).find(submission => submission.messageId === item.messageId); |
| 53 | } |
| 54 | |
| 55 | /** Pull same-turn output back behind its user when a published batch listed it first. */ |
| 56 | function userBeforeSameTurnOutput(items: readonly Item[]): Item[] { |
| 57 | const next = [...items]; |
| 58 | for (let userIndex = next.length - 1; userIndex >= 0; userIndex -= 1) { |
| 59 | const user = next[userIndex]; |
| 60 | if (user.kind !== "user" || !user.turnId) continue; |
| 61 | const ownedIds = new Set<string>(); |
| 62 | for (let index = 0; index < userIndex; index += 1) { |
| 63 | const item = next[index]; |
| 64 | if (item.kind !== "user" && item.turnId === user.turnId) ownedIds.add(item.id); |
| 65 | } |
| 66 | if (ownedIds.size === 0) continue; |
| 67 | const owned = next.filter(item => ownedIds.has(item.id)); |
| 68 | const rest = next.filter(item => !ownedIds.has(item.id)); |
| 69 | const at = rest.findIndex(item => item.id === user.id); |
| 70 | next.splice(0, next.length, ...rest.slice(0, at + 1), ...owned, ...rest.slice(at + 1)); |
| 71 | } |
| 72 | return next; |
| 73 | } |
| 74 | |
| 75 | export function installTranscriptRecords(s: State, a: Extract<Action, { type: "transcript_records" }>): State { |
| 76 | const ids = new Set<string>(); |
| 77 | for (const item of a.projection.items) { |
| 78 | if (!ids.has(item.id)) { ids.add(item.id); continue; } |
| 79 | recordFrontendDiagnostic("transcript", "projection.duplicate-item-id", { |
| 80 | itemId: item.id, projectedItems: a.projection.items.length, |
| 81 | }); |
| 82 | return s; |
| 83 | } |
| 84 | const projected = userBeforeSameTurnOutput(a.projection.items); |
| 85 | const projectedIds = new Set(projected.map(item => item.id)); |
| 86 | const removed = new Set(a.projection.removeIds); |
| 87 | const current = new Map<string, (typeof s.items)[number]>(); |
| 88 | for (const item of s.items) if (!current.has(item.id)) current.set(item.id, item); |
| 89 | const items = projected.map(update => { |
| 90 | const item = current.get(update.id); |
| 91 | if (!item) return update; |
| 92 | let merged = update; |
| 93 | if (item.kind === "tool" && update.kind === "tool") { |
| 94 | const runtime = { |
| 95 | parentId: update.parentId ?? item.parentId, |
| 96 | profile: update.profile ?? item.profile, |
| 97 | subagentProgress: update.subagentProgress ?? item.subagentProgress, |
| 98 | subagentOutcome: update.subagentOutcome ?? item.subagentOutcome, |
| 99 | verifying: update.verifying ?? item.verifying, |
| 100 | argChars: update.argChars ?? item.argChars, |
| 101 | }; |
| 102 | if (update.resultEvidence === "missing" || update.resultMissing) { |
| 103 | merged = { ...update, ...runtime, status: item.status, output: item.output, error: item.error, |
| 104 | execution: item.execution, presentedFiles: item.presentedFiles, durationMs: item.durationMs, |
| 105 | startedAt: item.startedAt, truncated: item.truncated }; |
| 106 | } else if (update.resultEvidence === "observation") { |
| 107 | merged = { ...update, ...runtime, output: item.output, error: item.error, |
| 108 | execution: update.execution ?? item.execution, presentedFiles: update.presentedFiles ?? item.presentedFiles, |
| 109 | durationMs: update.durationMs ?? item.durationMs, startedAt: item.startedAt, truncated: item.truncated }; |
| 110 | } else merged = { ...update, ...runtime }; |
| 111 | } else if (item.kind === "assistant" && update.kind === "assistant" && item.turnFinal && !update.turnFinal) { |
| 112 | merged = { ...update, turnFinal: true, turnDurationMs: item.turnDurationMs, turnUsage: item.turnUsage, |
| 113 | samplingCount: item.samplingCount, toolCount: item.toolCount }; |
| 114 | } |
| 115 | merged = retainResidentTurnId(merged, item); |
| 116 | return JSON.stringify(item) === JSON.stringify(merged) ? item : merged; |
| 117 | }); |
| 118 | // Preserve local/live rows at their nearest persisted anchors. The previous |
| 119 | // projection set distinguishes reclaimed formal rows from genuinely local |
| 120 | // rows, so a page removal cannot resurrect a stale transcript item. |
| 121 | const previousProjected = new Set(s.transcriptProjectedIds); |
| 122 | if (previousProjected.size === 0) { |
| 123 | for (const item of s.items) if (projectedIds.has(item.id)) previousProjected.add(item.id); |
| 124 | } |
| 125 | const localRows: Array<{ item: (typeof s.items)[number]; index: number; left?: string; right?: string }> = []; |
| 126 | for (let index = 0; index < s.items.length; index += 1) { |
| 127 | const item = s.items[index]; |
| 128 | if (previousProjected.has(item.id) || projectedIds.has(item.id) || removed.has(item.id)) continue; |
| 129 | let left: string | undefined; |
| 130 | let right: string | undefined; |
| 131 | for (let cursor = index - 1; cursor >= 0; cursor -= 1) { |
| 132 | if (projectedIds.has(s.items[cursor].id)) { left = s.items[cursor].id; break; } |
| 133 | } |
| 134 | for (let cursor = index + 1; cursor < s.items.length; cursor += 1) { |
| 135 | if (projectedIds.has(s.items[cursor].id)) { right = s.items[cursor].id; break; } |
| 136 | } |
| 137 | localRows.push({ item, index, left, right }); |
| 138 | } |
| 139 | // A durable user record can replace a local submission between two |
| 140 | // projections. Live process/assistant rows from that turn were previously |
| 141 | // anchored to the formal row before the local echo because the echo itself |
| 142 | // is owned outside `items`. Move that anchor to the durable user record so |
| 143 | // the already-mounted turn stays in user -> process -> assistant order. |
| 144 | const turnAnchors = new Map<string, string>(); |
| 145 | const userOwners = new Map<string, string | undefined>(); |
| 146 | let userId: string | undefined; |
| 147 | for (const item of items) { |
| 148 | if (item.kind === "user") userId = item.id; |
| 149 | userOwners.set(item.id, userId); |
| 150 | if (item.kind === "user" && item.turnId) turnAnchors.set(item.turnId, item.id); |
| 151 | } |
| 152 | for (const local of Object.values(s.localSubmissions)) { |
| 153 | const id = local.messageId ? `m:${local.messageId}` : undefined; |
| 154 | if (local.turnId && id && projectedIds.has(id)) turnAnchors.set(local.turnId, id); |
| 155 | } |
| 156 | // A published user row can arrive without turnId. The send-time anchor still |
| 157 | // owns every live row that appeared after that send. |
| 158 | const submissionTails = new Map<string, string>(); |
| 159 | const anchoredSubmissions = new Set<string>(); |
| 160 | for (const item of projected) { |
| 161 | if (item.kind !== "user") continue; |
| 162 | const submission = submissionForInstalledUser(s, item); |
| 163 | if (!submission || submission.placement === "latest" || anchoredSubmissions.has(submission.submissionId)) continue; |
| 164 | anchoredSubmissions.add(submission.submissionId); |
| 165 | submissionTails.set(submission.submissionId, item.id); |
| 166 | } |
| 167 | const tailAnchor = (row: (typeof localRows)[number]): string | undefined => { |
| 168 | let owner: string | undefined; |
| 169 | for (const submission of orderedLocalSubmissions(s)) { |
| 170 | const userId = submissionTails.get(submission.submissionId); |
| 171 | if (!userId) continue; |
| 172 | if (!submission.anchorItemId && submission.placement !== "after") { |
| 173 | if (!s.historyHasOlder && s.historyStartTurn === 0) owner = userId; |
| 174 | continue; |
| 175 | } |
| 176 | if (!submission.anchorItemId) continue; |
| 177 | const anchorIndex = s.items.findIndex(candidate => candidate.id === submission.anchorItemId); |
| 178 | if (anchorIndex >= 0 && row.index > anchorIndex) owner = userId; |
| 179 | } |
| 180 | return owner; |
| 181 | }; |
| 182 | const tails = new Map<string, string>(); |
| 183 | for (const row of localRows) { |
| 184 | const owner = (row.item.turnId && turnAnchors.get(row.item.turnId)) || tailAnchor(row); |
| 185 | // The user bounds the turn; a surviving (or newly formal) predecessor |
| 186 | // within that turn still owns this row's position among samples and tools. |
| 187 | const anchor = owner && (!row.left || userOwners.get(row.left) !== owner) ? owner : row.left; |
| 188 | const left = anchor && (tails.get(anchor) ?? anchor); |
| 189 | const leftIndex = left ? items.findIndex(item => item.id === left) : -1; |
| 190 | if (leftIndex >= 0) { |
| 191 | items.splice(leftIndex + 1, 0, row.item); |
| 192 | if (anchor) tails.set(anchor, row.item.id); |
| 193 | continue; |
| 194 | } |
| 195 | const rightIndex = row.right ? items.findIndex(item => item.id === row.right) : -1; |
| 196 | if (rightIndex >= 0) items.splice(rightIndex, 0, row.item); |
| 197 | else items.push(row.item); |
| 198 | } |
| 199 | const mutation = a.projection.mutation ?? "patch"; |
| 200 | const projectionIds = [...projectedIds]; |
| 201 | const structural = items.length !== s.items.length || items.some((item, index) => item !== s.items[index]) |
| 202 | || projectionIds.length !== s.transcriptProjectedIds.length |
| 203 | || projectionIds.some((id, index) => id !== s.transcriptProjectedIds[index]); |
| 204 | const digest = a.projection.digest || s.historyDigest; |
| 205 | const revision = a.projection.revisionKnown ? a.projection.revision : s.historyRevision; |
| 206 | const metadata = s.historyPrefixCount !== projected.length || s.historyStartTurn !== a.projection.startTurn |
| 207 | || s.historyEndTurn !== a.projection.endTurn || s.historyTotalTurns !== a.projection.totalTurns |
| 208 | || s.historyHasOlder !== a.projection.hasOlder || s.historyHasNewer !== a.projection.hasNewer |
| 209 | || s.historyOlderLoading || s.historyNewerLoading || s.historyOlderError !== undefined |
| 210 | || s.historyNewerError !== undefined || s.historyRevision !== revision || s.historyDigest !== digest; |
| 211 | const confirmations = [...a.confirmedUsers, ...canonicalUserConfirmations(a.projection.items)]; |
| 212 | const settled = settleLocalSubmissions(s, items, confirmations); |
| 213 | if (!structural && !metadata && settled === s) return s; |
| 214 | return { ...settled, items, transcriptProjectedIds: projectionIds, |
| 215 | historyPrefixCount: projected.length, |
| 216 | historyStartTurn: a.projection.startTurn, historyEndTurn: a.projection.endTurn, |
| 217 | historyTotalTurns: a.projection.totalTurns, |
| 218 | historyHasOlder: a.projection.hasOlder, historyHasNewer: a.projection.hasNewer, |
| 219 | historyOlderLoading: false, historyOlderError: undefined, |
| 220 | historyNewerLoading: false, historyNewerError: undefined, |
| 221 | historyRevision: revision, historyDigest: digest, |
| 222 | historyLayoutRevision: s.historyLayoutRevision + (structural ? 1 : 0), |
| 223 | historyMutation: structural |
| 224 | ? { seq: s.historyMutation.seq + 1, kind: mutation === "patch" ? "patch" : mutation } |
| 225 | : s.historyMutation }; |
| 226 | } |
| 227 | |
| 228 | export function startLocalSubmission(s: State, a: Extract<Action, { type: "user" }>, clock: number): State { |
| 229 | const seq = a.seq !== undefined ? a.seq : s.seq; |
| 230 | const userItemId = `u${seq}`; |
| 231 | const next = { |
| 232 | ...s, |
| 233 | completionSummary: undefined, |
| 234 | seq: seq + 1, |
| 235 | items: s.items.map(item => item.kind==="notice" && item.action==="recover_context" ? {...item,action:undefined} : item), |
| 236 | running: true, |
| 237 | pendingPrompt: false, |
| 238 | cancelRequested: false, |
| 239 | cancellable: true, |
| 240 | ...resetTurnTiming(), |
| 241 | turnLifecycleObservedAt: clock, |
| 242 | // New turn epoch: forget the previous prompt anchor so a genuinely new |
| 243 | // prompt re-anchors freshly instead of inheriting a stale id/time. |
| 244 | promptArrivedAt: undefined, |
| 245 | promptArrivedId: undefined, |
| 246 | pendingUser: a.text, |
| 247 | pendingSubmissionId: a.submissionId, |
| 248 | activeTurnId: s.turnActive ? s.activeTurnId : undefined, |
| 249 | currentAssistant: undefined, |
| 250 | assistantSegmentOrdinal: 0, |
| 251 | live: undefined, |
| 252 | streamAttemptJournal: undefined, |
| 253 | streamInterruptNoticeShown: undefined, |
| 254 | deliveryRecoveryActive: Boolean(a.deliveryRecovery), |
| 255 | discardTurn: false, |
| 256 | }; |
| 257 | return beginLocalSubmission(next, { |
| 258 | submissionId: a.submissionId, |
| 259 | localId: userItemId, |
| 260 | text: a.text, |
| 261 | submitText: a.submitText, |
| 262 | createdAt: Date.now(), |
| 263 | sequence: seq, |
| 264 | anchorItemId: s.historyHasNewer ? undefined : s.items[s.items.length - 1]?.id, |
| 265 | placement: s.historyHasNewer ? "latest" : s.items.length ? "after" : "start", |
| 266 | }); |
| 267 | } |
| 268 | |
| 269 | /** Stamp this event's turn onto rows that just arrived or still belong to the live answer. */ |
| 270 | export function stampArrivingTurnId<T extends { items: Item[]; currentAssistant?: string; activeTurnId?: string }>(state: T, before: readonly Item[], turnId: string | undefined, messageId?: string): T { |
| 271 | if (!turnId) return state; |
| 272 | const foreign = Boolean(state.activeTurnId) && state.activeTurnId !== turnId; |
| 273 | const prior = new Set(before); |
| 274 | let stamped = false; |
| 275 | const items = state.items.map(item => { |
| 276 | if (item.turnId) return item; |
| 277 | const createdNow = !prior.has(item); |
| 278 | const liveAssistant = !foreign && item.kind === "assistant" && item.id === state.currentAssistant; |
| 279 | const liveTool = !foreign && item.kind === "tool" && Boolean(messageId) && item.messageId === messageId; |
| 280 | if (!createdNow && !liveAssistant && !liveTool) return item; |
| 281 | stamped = true; |
| 282 | return { ...item, turnId }; |
| 283 | }); |
| 284 | return stamped ? { ...state, items } : state; |
| 285 | } |
| 286 |