| 1 | import type { HistoryMessage, WireSessionOperation } from "./types"; |
| 2 | import type { Item } from "./useController"; |
| 3 | import type { StructuredInvocationSubmit } from "./invocationDisplay"; |
| 4 | |
| 5 | export type CompactionItem = Extract<Item, { kind: "compaction" }>; |
| 6 | |
| 7 | /** Match the controller's management route without treating other slash inputs as commands. */ |
| 8 | export function isCompactCommand(input: string): boolean { |
| 9 | const trimmed = input.trim(); |
| 10 | return trimmed === "/compact" || trimmed.startsWith("/compact "); |
| 11 | } |
| 12 | |
| 13 | /** Structured invocations and initial goals retain their own admission contracts. */ |
| 14 | export function isCompactSubmission(input: string, structured?: StructuredInvocationSubmit, initialGoal?: object): boolean { |
| 15 | return !structured && !initialGoal && isCompactCommand(input); |
| 16 | } |
| 17 | |
| 18 | const terminalStatuses = new Set([ |
| 19 | "completed", |
| 20 | "cancelled", |
| 21 | "partially_completed", |
| 22 | "failed", |
| 23 | "noop", |
| 24 | "interrupted", |
| 25 | "recovery_required", |
| 26 | ]); |
| 27 | const pendingStatuses = new Set(["running", "cancelling", "finalizing", "loading", "confirming"]); |
| 28 | |
| 29 | function object(value: unknown): Record<string, unknown> | undefined { |
| 30 | return value && typeof value === "object" && !Array.isArray(value) ? value as Record<string, unknown> : undefined; |
| 31 | } |
| 32 | |
| 33 | function text(value: unknown): string | undefined { |
| 34 | return typeof value === "string" && value !== "" ? value : undefined; |
| 35 | } |
| 36 | |
| 37 | function count(value: unknown): number | undefined { |
| 38 | return typeof value === "number" && Number.isSafeInteger(value) && value >= 0 ? value : undefined; |
| 39 | } |
| 40 | |
| 41 | function boolean(value: unknown): boolean | undefined { |
| 42 | return typeof value === "boolean" ? value : undefined; |
| 43 | } |
| 44 | |
| 45 | /** Decodes the content of a session-maintenance-v1 compaction display row. */ |
| 46 | export function parseSessionOperation(value: unknown): WireSessionOperation | undefined { |
| 47 | let raw = object(value); |
| 48 | if (!raw) return undefined; |
| 49 | if (!text(raw.operationId) && typeof raw.content === "string") { |
| 50 | try { raw = object(JSON.parse(raw.content)) ?? raw; } catch { return undefined; } |
| 51 | } |
| 52 | const operationId = text(raw.operationId); |
| 53 | if (!operationId) return undefined; |
| 54 | return { |
| 55 | operationId, |
| 56 | kind: text(raw.kind) ?? "compact", |
| 57 | activity: text(raw.activity) ?? "", |
| 58 | status: text(raw.status) ?? "", |
| 59 | errorCode: text(raw.errorCode), |
| 60 | detail: text(raw.detail), |
| 61 | applied: boolean(raw.applied), |
| 62 | inputTokens: count(raw.inputTokens), |
| 63 | resultTokens: count(raw.resultTokens), |
| 64 | messages: count(raw.messages), |
| 65 | summary: typeof raw.summary === "string" ? raw.summary : undefined, |
| 66 | archive: typeof raw.archive === "string" ? raw.archive : undefined, |
| 67 | operationRevision: count(raw.operationRevision), |
| 68 | runtimeEpoch: text(raw.runtimeEpoch), |
| 69 | }; |
| 70 | } |
| 71 | |
| 72 | export function isTerminalSessionOperation(status: string | undefined): boolean { |
| 73 | return Boolean(status && terminalStatuses.has(status)); |
| 74 | } |
| 75 | |
| 76 | export function sessionOperationStatus(status?: string, activity?: string): string { |
| 77 | const value = status || activity || "unavailable"; |
| 78 | return isTerminalSessionOperation(value) || pendingStatuses.has(value) |
| 79 | ? value : "unavailable"; |
| 80 | } |
| 81 | |
| 82 | function operationPending(status: string | undefined): boolean { |
| 83 | return pendingStatuses.has(status ?? ""); |
| 84 | } |
| 85 | |
| 86 | export function sessionOperationFromHistory(message: HistoryMessage): WireSessionOperation | undefined { |
| 87 | if (message.role !== "compaction" || !message.operationId) return undefined; |
| 88 | const status = sessionOperationStatus(message.operationStatus, message.operationActivity); |
| 89 | const awaitingRuntime = Boolean(message.pending && operationPending(status) && status !== "loading"); |
| 90 | return { |
| 91 | operationId: message.operationId, |
| 92 | kind: message.operationKind ?? "compact", |
| 93 | activity: message.operationActivity ?? "", |
| 94 | status: awaitingRuntime ? "confirming" : status, |
| 95 | errorCode: message.errorCode, |
| 96 | detail: message.detail, |
| 97 | applied: message.applied, |
| 98 | inputTokens: message.inputTokens, |
| 99 | resultTokens: message.resultTokens, |
| 100 | messages: message.messages, |
| 101 | summary: message.summary, |
| 102 | archive: message.archive, |
| 103 | operationRevision: message.operationRevision, |
| 104 | runtimeEpoch: message.runtimeEpoch, |
| 105 | }; |
| 106 | } |
| 107 | |
| 108 | export function sessionOperationHistoryMessage(messageId: string, operation: WireSessionOperation): HistoryMessage { |
| 109 | const status = sessionOperationStatus(operation.status, operation.activity); |
| 110 | return { |
| 111 | role: "compaction", |
| 112 | messageId, |
| 113 | content: "", |
| 114 | pending: operationPending(status), |
| 115 | trigger: "manual", |
| 116 | messages: operation.messages, |
| 117 | summary: operation.summary, |
| 118 | archive: operation.archive, |
| 119 | operationId: operation.operationId, |
| 120 | operationKind: operation.kind, |
| 121 | operationStatus: status, |
| 122 | operationActivity: operation.activity || undefined, |
| 123 | operationRevision: operation.operationRevision, |
| 124 | runtimeEpoch: operation.runtimeEpoch, |
| 125 | errorCode: operation.errorCode, |
| 126 | detail: operation.detail, |
| 127 | applied: operation.applied, |
| 128 | inputTokens: operation.inputTokens, |
| 129 | resultTokens: operation.resultTokens, |
| 130 | }; |
| 131 | } |
| 132 | |
| 133 | export function sessionOperationItem(operation: WireSessionOperation, historyEntryId?: string): CompactionItem { |
| 134 | const status = sessionOperationStatus(operation.status, operation.activity); |
| 135 | return { |
| 136 | ...operation, |
| 137 | kind: "compaction", |
| 138 | id: `maintenance:${operation.operationId}`, |
| 139 | pending: operationPending(status), |
| 140 | trigger: "manual", |
| 141 | messages: operation.messages ?? 0, |
| 142 | summary: operation.summary ?? "", |
| 143 | archive: operation.archive ?? "", |
| 144 | operationKind: operation.kind, |
| 145 | status, |
| 146 | activity: operation.activity || undefined, |
| 147 | historyEntryId, |
| 148 | }; |
| 149 | } |
| 150 | |
| 151 | function fillMissing(prior: CompactionItem, incoming: CompactionItem): CompactionItem { |
| 152 | return { |
| 153 | ...prior, |
| 154 | operationKind: prior.operationKind || incoming.operationKind, |
| 155 | errorCode: prior.errorCode ?? incoming.errorCode, |
| 156 | detail: prior.detail ?? incoming.detail, |
| 157 | applied: prior.applied || incoming.applied || undefined, |
| 158 | inputTokens: prior.inputTokens ?? incoming.inputTokens, |
| 159 | resultTokens: prior.resultTokens ?? incoming.resultTokens, |
| 160 | messages: prior.messages || incoming.messages, |
| 161 | summary: prior.summary || incoming.summary, |
| 162 | archive: prior.archive || incoming.archive, |
| 163 | operationRevision: prior.operationRevision ?? incoming.operationRevision, |
| 164 | runtimeEpoch: prior.runtimeEpoch ?? incoming.runtimeEpoch, |
| 165 | historyEntryId: prior.historyEntryId ?? incoming.historyEntryId, |
| 166 | }; |
| 167 | } |
| 168 | |
| 169 | /** |
| 170 | * Merges durable history, live events and runtime snapshots for one operation. |
| 171 | * Missing fields are partial observations; revisions and terminality are |
| 172 | * monotonic so a delayed progress update cannot reopen a completed card. |
| 173 | */ |
| 174 | export function mergeSessionOperationItem(prior: CompactionItem | undefined, incoming: CompactionItem): CompactionItem { |
| 175 | if (!prior || prior.operationId !== incoming.operationId) return incoming; |
| 176 | const priorRevision = prior.operationRevision; |
| 177 | const incomingRevision = incoming.operationRevision; |
| 178 | const older = priorRevision !== undefined && incomingRevision !== undefined && incomingRevision < priorRevision; |
| 179 | const epochMismatch = Boolean(prior.runtimeEpoch && incoming.runtimeEpoch && prior.runtimeEpoch !== incoming.runtimeEpoch); |
| 180 | const terminalRegression = !prior.interruptionInferred && isTerminalSessionOperation(prior.status) && !isTerminalSessionOperation(incoming.status); |
| 181 | if (older || epochMismatch || terminalRegression) return fillMissing(prior, incoming); |
| 182 | |
| 183 | const merged: CompactionItem = { |
| 184 | ...prior, |
| 185 | ...Object.fromEntries(Object.entries(incoming).filter(([, value]) => value !== undefined)), |
| 186 | id: prior.id, |
| 187 | messages: incoming.messages || prior.messages, |
| 188 | summary: incoming.summary || prior.summary, |
| 189 | archive: incoming.archive || prior.archive, |
| 190 | applied: prior.applied || incoming.applied || undefined, |
| 191 | interruptionInferred: incoming.interruptionInferred, |
| 192 | } as CompactionItem; |
| 193 | merged.pending = operationPending(merged.status); |
| 194 | return merged; |
| 195 | } |
| 196 | |
| 197 | /** Upserts by operation identity and collapses any legacy duplicate cards. */ |
| 198 | export function upsertSessionOperationItem(items: Item[], incoming: CompactionItem): { items: Item[]; inserted: boolean } { |
| 199 | const indexes: number[] = []; |
| 200 | items.forEach((item, index) => { |
| 201 | if (item.kind === "compaction" && item.operationId === incoming.operationId) indexes.push(index); |
| 202 | }); |
| 203 | if (indexes.length === 0) return { items: [...items, incoming], inserted: true }; |
| 204 | let merged = incoming; |
| 205 | for (const index of indexes) merged = mergeSessionOperationItem(items[index] as CompactionItem, merged); |
| 206 | const first = indexes[0]; |
| 207 | return { |
| 208 | items: items.flatMap((item, index) => index === first ? [merged] : indexes.includes(index) ? [] : [item]), |
| 209 | inserted: false, |
| 210 | }; |
| 211 | } |
| 212 | |
| 213 | /** |
| 214 | * Reconciles a freshly loaded history window with operation updates that may |
| 215 | * have arrived while the read was in flight. Only matching identities merge; |
| 216 | * an active card absent from the window is retained as a live tail. |
| 217 | */ |
| 218 | export function reconcileSessionOperationItems(history: Item[], live: Item[], preserveMissing = false): Item[] { |
| 219 | let result = history; |
| 220 | for (const item of live) { |
| 221 | if (item.kind !== "compaction" || !item.operationId) continue; |
| 222 | if (result.some(candidate => candidate.kind === "compaction" && candidate.operationId === item.operationId)) { |
| 223 | result = upsertSessionOperationItem(result, item).items; |
| 224 | } else if (preserveMissing && item.pending) { |
| 225 | result = [...result, item]; |
| 226 | } |
| 227 | } |
| 228 | return result; |
| 229 | } |
| 230 | |
| 231 | /** |
| 232 | * Resolves durable progress rows once a runtime snapshot has established the |
| 233 | * active maintenance identity. A history row remains pending until that first |
| 234 | * snapshot arrives; afterwards only the matching runtime operation may remain |
| 235 | * pending. This is deliberately a display-only recovery decision. |
| 236 | */ |
| 237 | export function interruptOrphanedSessionOperationItems( |
| 238 | items: Item[], |
| 239 | activeOperationId?: string, |
| 240 | activeRuntimeEpoch?: string, |
| 241 | runtimeRevision?: number, |
| 242 | ): Item[] { |
| 243 | let changed = false; |
| 244 | const next = items.map((item) => { |
| 245 | if (item.kind !== "compaction" || !item.operationId || !item.pending || item.status === "loading" || isTerminalSessionOperation(item.status)) return item; |
| 246 | // A snapshot cached before this live event cannot declare it orphaned. |
| 247 | // A changed epoch is independently authoritative even if revisions reset. |
| 248 | if ((!item.runtimeEpoch || item.runtimeEpoch === activeRuntimeEpoch) && runtimeRevision !== undefined |
| 249 | && item.observedRuntimeRevision !== undefined && runtimeRevision <= item.observedRuntimeRevision) return item; |
| 250 | const identityMatches = item.operationId === activeOperationId |
| 251 | && !(item.runtimeEpoch && activeRuntimeEpoch && item.runtimeEpoch !== activeRuntimeEpoch); |
| 252 | if (identityMatches) return item; |
| 253 | changed = true; |
| 254 | return { ...item, pending: false, status: "interrupted", activity: "interrupted", interruptionInferred: true }; |
| 255 | }); |
| 256 | return changed ? next : items; |
| 257 | } |
| 258 |