| 1 | import type { Action, Item, State } from "./useController"; |
| 2 | import type { WireEvent } from "./types"; |
| 3 | import { acceptSessionRuntimeSnapshot, type RuntimeState } from "./runtimeStateStore"; |
| 4 | import { interruptOrphanedSessionOperationItems, sessionOperationItem, upsertSessionOperationItem } from "./sessionMaintenanceOperation"; |
| 5 | |
| 6 | function lastPendingCompaction(items: readonly Item[], manual = false): number { |
| 7 | for (let i = items.length - 1; i >= 0; i--) { |
| 8 | const item = items[i]; |
| 9 | if (item.kind === "compaction" && item.pending && (!manual || item.operationId)) return i; |
| 10 | } |
| 11 | return -1; |
| 12 | } |
| 13 | |
| 14 | export function reduceCompactionEvent(s: State, e: WireEvent): State { |
| 15 | switch (e.kind) { |
| 16 | case "session_operation": { |
| 17 | const op = e.sessionOperation; |
| 18 | if (!op?.operationId) return s; |
| 19 | const updated = upsertSessionOperationItem(s.items, { ...sessionOperationItem(op), observedRuntimeRevision: s.runtimeStateSnapshot?.revision ?? 0 }); |
| 20 | return { ...s, seq: s.seq + (updated.inserted ? 1 : 0), items: updated.items }; |
| 21 | } |
| 22 | case "compaction_started": |
| 23 | if (lastPendingCompaction(s.items, true) >= 0) return s; |
| 24 | return { ...s, seq: s.seq + 1, items: [...s.items, { kind: "compaction", id: `c${s.seq}`, pending: true, trigger: e.compaction?.trigger ?? "", messages: 0, summary: "", archive: "" }] }; |
| 25 | case "compaction_done": { |
| 26 | const c = e.compaction; |
| 27 | const operationAt = lastPendingCompaction(s.items, true); |
| 28 | if (operationAt >= 0) { |
| 29 | const items = s.items.map((it, i) => i === operationAt && it.kind === "compaction" ? { ...it, |
| 30 | messages: c?.messages ?? it.messages, summary: c?.summary ?? it.summary, archive: c?.archive ?? it.archive } : it); |
| 31 | return { ...s, items }; |
| 32 | } |
| 33 | const at = lastPendingCompaction(s.items); |
| 34 | if (!c?.summary) { |
| 35 | const items = at < 0 ? s.items : s.items.filter((_, i) => i !== at); |
| 36 | return { ...s, running: s.turnActive ? s.running : false, items }; |
| 37 | } |
| 38 | const filled: Item = { kind: "compaction", id: at < 0 ? `c${s.seq}` : (s.items[at] as Extract<Item, { kind: "compaction" }>).id, pending: false, trigger: c.trigger ?? "", messages: c.messages ?? 0, summary: c.summary, archive: c.archive ?? "" }; |
| 39 | const items = at < 0 ? [...s.items, filled] : s.items.map((it, i) => (i === at ? filled : it)); |
| 40 | return { ...s, running: s.turnActive ? s.running : false, seq: s.seq + 1, items }; |
| 41 | } |
| 42 | default: return s; |
| 43 | } |
| 44 | } |
| 45 | |
| 46 | export function reduceMaintenanceRuntimeSnapshot(s: State, snapshot: RuntimeState): State { |
| 47 | const runtimeStateSnapshot = acceptSessionRuntimeSnapshot(s.runtimeStateSnapshot, snapshot); |
| 48 | if (runtimeStateSnapshot === s.runtimeStateSnapshot) return s; |
| 49 | const meta = runtimeStateSnapshot.todos !== undefined && s.meta |
| 50 | ? { ...s.meta, runtimeStateSnapshot, canonicalTodos: runtimeStateSnapshot.todos } |
| 51 | : s.meta; |
| 52 | const maintenance = runtimeStateSnapshot.maintenance; |
| 53 | if (!maintenance) return { ...s, meta, runtimeStateSnapshot }; |
| 54 | const status = ["cancelling", "finalizing", "recovery_required"].includes(maintenance.activity) |
| 55 | ? maintenance.activity : "running"; |
| 56 | const updated = upsertSessionOperationItem(s.items, { ...sessionOperationItem({ |
| 57 | ...maintenance, |
| 58 | status: maintenance.status || status, |
| 59 | runtimeEpoch: maintenance.runtimeEpoch || runtimeStateSnapshot.runtimeEpoch, |
| 60 | }), observedRuntimeRevision: runtimeStateSnapshot.revision }); |
| 61 | return { ...s, meta, runtimeStateSnapshot, items: updated.items, seq: s.seq + (updated.inserted ? 1 : 0) }; |
| 62 | } |
| 63 | |
| 64 | export function reconcileMaintenanceState(next: State, a: Action): State { |
| 65 | const historyOrRuntimeSynchronized = ["runtime_snapshot", "meta", "history", "history_page", |
| 66 | "history_replace", "history_rebase", "history_prepend", "history_append", "history_items_patch", |
| 67 | "transcript_snapshot", "transcript_v2_snapshot", "transcript_page", "transcript_records"].includes(a.type); |
| 68 | if (historyOrRuntimeSynchronized && next.runtimeStateSnapshot) { |
| 69 | const maintenance = next.runtimeStateSnapshot.maintenance; |
| 70 | const items = interruptOrphanedSessionOperationItems( |
| 71 | next.items, |
| 72 | maintenance?.operationId, |
| 73 | maintenance?.runtimeEpoch || next.runtimeStateSnapshot.runtimeEpoch, |
| 74 | next.runtimeStateSnapshot.revision, |
| 75 | ); |
| 76 | if (items !== next.items) next = { ...next, items }; |
| 77 | } |
| 78 | return next; |
| 79 | } |
| 80 |