| 1 | import { PET_MAX_SECONDS } from './pet-sim.js'; |
| 2 | import type { Category, WhaleEvent } from './model.js'; |
| 3 | import { compilePetTelemetry, PET_BIN_MS, type PetBucket } from './pet-telemetry.js'; |
| 4 | |
| 5 | // Engine emits a liveness pulse every 10 seconds for a running operation. |
| 6 | // Keep a small scheduling margin so one delayed pulse does not erase a valid |
| 7 | // long operation from the owner projection; after that, it becomes generic. |
| 8 | // The physics tape keeps its own, shorter coverage rule (see `pulse`). |
| 9 | export const ENGINE_OWNER_STALE_MS = 12_000; |
| 10 | export type OwnerActivityKind = 'reading' | 'editing' | 'searching' | 'testing' | 'executing' |
| 11 | | 'browsing' | 'computer' | 'memory' | 'tool' | 'thinking' | 'responding' | 'delegating'; |
| 12 | export type OwnerPresence = 'unknown' | 'working' | 'needs_you' | 'done' | 'idle'; |
| 13 | export type OwnerFreshness = 'missing' | 'fresh' | 'stale'; |
| 14 | export type OwnerTurnOutcome = 'completed' | 'interrupted' | 'failed'; |
| 15 | export type OwnerOperationOutcome = 'succeeded' | 'failed' | 'cancelled' | 'denied'; |
| 16 | |
| 17 | export interface EngineOwnerProjection { |
| 18 | schemaVersion: 1; |
| 19 | /** Added by the existing owner after selecting its session source. */ |
| 20 | sessionId: string | null; |
| 21 | /** Added by the existing owner; monotonic for this session source. */ |
| 22 | cursor: number; |
| 23 | observed: boolean; |
| 24 | freshness: OwnerFreshness; |
| 25 | authoritativePresence: OwnerPresence; |
| 26 | activityKind: OwnerActivityKind | null; |
| 27 | observedAtMs: number | null; |
| 28 | parallelAgentCount: number; |
| 29 | activeSpans: { activityKind: OwnerActivityKind; startedAtMs: number }[]; |
| 30 | turnId: string | null; |
| 31 | turnOutcome: OwnerTurnOutcome | null; |
| 32 | doneEffectId: string | null; |
| 33 | failedToolAge: { activityKind: OwnerActivityKind; ageMs: number } | null; |
| 34 | } |
| 35 | |
| 36 | type SpanGroup = 'operation' | 'agent' | 'thinking' | 'responding'; |
| 37 | interface ActiveSpan { |
| 38 | key: string; |
| 39 | activityKind: OwnerActivityKind; |
| 40 | startedAtMs: number; |
| 41 | event: WhaleEvent; |
| 42 | group: SpanGroup; |
| 43 | } |
| 44 | interface FailedTool { |
| 45 | activityKind: OwnerActivityKind; |
| 46 | failedAtMs: number; |
| 47 | } |
| 48 | |
| 49 | const ACTIVITY_KINDS: readonly OwnerActivityKind[] = [ |
| 50 | 'reading', 'editing', 'searching', 'testing', 'executing', 'browsing', 'computer', |
| 51 | 'memory', 'tool', 'thinking', 'responding', 'delegating', |
| 52 | ]; |
| 53 | const TURN_OUTCOMES: readonly OwnerTurnOutcome[] = ['completed', 'interrupted', 'failed']; |
| 54 | const OPERATION_OUTCOMES: readonly OwnerOperationOutcome[] = ['succeeded', 'failed', 'cancelled', 'denied']; |
| 55 | const APPROVAL_OUTCOMES = ['approved', 'denied', 'cancelled'] as const; |
| 56 | |
| 57 | function categoryFor(kind: OwnerActivityKind): Category { |
| 58 | switch (kind) { |
| 59 | case 'reading': case 'editing': case 'searching': return 'filesystem'; |
| 60 | case 'testing': case 'executing': return 'code'; |
| 61 | case 'browsing': return 'browser'; |
| 62 | case 'computer': case 'tool': return 'tool'; |
| 63 | case 'memory': return 'memory'; |
| 64 | case 'thinking': return 'reasoning'; |
| 65 | case 'responding': return 'communication'; |
| 66 | case 'delegating': return 'agent'; |
| 67 | } |
| 68 | } |
| 69 | |
| 70 | /** Read-only reducer for the Engine owner's typed metadata. It keeps using the |
| 71 | * existing pet telemetry tape and physics owner; it neither recognizes tool |
| 72 | * names nor accepts transcript text, arguments, commands, or results. */ |
| 73 | export class PetEngineTelemetry { |
| 74 | private events: WhaleEvent[] = []; |
| 75 | private active = new Map<string, ActiveSpan>(); |
| 76 | private waiting: WhaleEvent | undefined; |
| 77 | private sequence = 0; |
| 78 | private lastTime = 0; |
| 79 | private lastObservedAt: number | undefined; |
| 80 | private turnId: string | undefined; |
| 81 | private turnOutcome: OwnerTurnOutcome | undefined; |
| 82 | private terminalAt: number | undefined; |
| 83 | private lastFailedTool: FailedTool | undefined; |
| 84 | private completedSpans = new Set<string>(); |
| 85 | private completedTurns = new Set<string>(); |
| 86 | |
| 87 | /** Ephemeral safe read projection. Span ids are retained only in this |
| 88 | * reducer to correlate trusted lifecycle events and never leave the owner. */ |
| 89 | activity(at: number): EngineOwnerProjection { |
| 90 | const observed = this.lastObservedAt !== undefined; |
| 91 | const freshness: OwnerFreshness = !observed ? 'missing' |
| 92 | : at - this.lastObservedAt! <= ENGINE_OWNER_STALE_MS ? 'fresh' : 'stale'; |
| 93 | const fresh = freshness === 'fresh'; |
| 94 | const active = fresh |
| 95 | ? [...this.active.values()] |
| 96 | .filter(span => at >= span.startedAtMs && at - span.event.endTime <= ENGINE_OWNER_STALE_MS) |
| 97 | // The parent's own work leads; delegated agents are counted in |
| 98 | // `parallelAgentCount` and only lead when nothing else is active. |
| 99 | .sort((a, b) => Number(a.group === 'agent') - Number(b.group === 'agent') |
| 100 | || b.startedAtMs - a.startedAtMs || a.key.localeCompare(b.key)) |
| 101 | : []; |
| 102 | const agents = active.filter(span => span.group === 'agent').length; |
| 103 | const terminalFresh = fresh && this.terminalAt !== undefined |
| 104 | && at >= this.terminalAt && at - this.terminalAt <= ENGINE_OWNER_STALE_MS; |
| 105 | let authoritativePresence: OwnerPresence = 'unknown'; |
| 106 | if (fresh) { |
| 107 | if (this.waiting) authoritativePresence = 'needs_you'; |
| 108 | else if (terminalFresh && this.turnOutcome === 'completed' && this.turnId) authoritativePresence = 'done'; |
| 109 | else if (terminalFresh) authoritativePresence = 'idle'; |
| 110 | else if (active.length > 0 || this.turnId) authoritativePresence = 'working'; |
| 111 | } |
| 112 | const activityKind = fresh && authoritativePresence !== 'needs_you' |
| 113 | && authoritativePresence !== 'done' && authoritativePresence !== 'idle' |
| 114 | ? active[0]?.activityKind |
| 115 | ?? (this.lastFailedTool && at >= this.lastFailedTool.failedAtMs |
| 116 | && at - this.lastFailedTool.failedAtMs <= ENGINE_OWNER_STALE_MS |
| 117 | ? this.lastFailedTool.activityKind : null) |
| 118 | : null; |
| 119 | const failedAge = fresh && this.lastFailedTool |
| 120 | && at >= this.lastFailedTool.failedAtMs |
| 121 | && at - this.lastFailedTool.failedAtMs <= ENGINE_OWNER_STALE_MS |
| 122 | ? { activityKind: this.lastFailedTool.activityKind, ageMs: at - this.lastFailedTool.failedAtMs } |
| 123 | : null; |
| 124 | const doneEffectId = authoritativePresence === 'done' && terminalFresh |
| 125 | && this.turnOutcome === 'completed' && this.turnId ? this.turnId : null; |
| 126 | |
| 127 | return { |
| 128 | schemaVersion: 1, |
| 129 | sessionId: null, |
| 130 | cursor: 0, |
| 131 | observed, |
| 132 | freshness, |
| 133 | authoritativePresence, |
| 134 | activityKind: activityKind ?? null, |
| 135 | observedAtMs: this.lastObservedAt ?? null, |
| 136 | parallelAgentCount: fresh ? agents : 0, |
| 137 | activeSpans: active.slice(0, 4).map(span => ({ |
| 138 | activityKind: span.activityKind, |
| 139 | startedAtMs: span.startedAtMs, |
| 140 | })), |
| 141 | turnId: this.turnId ?? null, |
| 142 | turnOutcome: this.turnOutcome ?? null, |
| 143 | doneEffectId, |
| 144 | failedToolAge: failedAge, |
| 145 | }; |
| 146 | } |
| 147 | |
| 148 | private add(name: string, category: Category, at: number, agentId = 'parent', continuation = false): WhaleEvent { |
| 149 | if (this.events.length >= 8192) throw new Error('Pet Engine observation window is full.'); |
| 150 | const event: WhaleEvent = { |
| 151 | schemaVersion: 1, |
| 152 | id: `engine:${this.sequence++}`, |
| 153 | traceId: 'foreground', |
| 154 | startTime: at, |
| 155 | endTime: at, |
| 156 | name, |
| 157 | category, |
| 158 | agentId, |
| 159 | status: 'running', |
| 160 | attributes: continuation ? { 'whalesong.continuation': true } : {}, |
| 161 | }; |
| 162 | this.events.push(event); |
| 163 | return event; |
| 164 | } |
| 165 | |
| 166 | private start(key: string, kind: OwnerActivityKind, group: SpanGroup, at: number, agentId = 'parent'): void { |
| 167 | if (this.active.has(key)) return; |
| 168 | if (this.active.size >= 256) throw new Error('Too many active Engine pet spans.'); |
| 169 | const event = this.add(group === 'agent' ? 'agent' : group === 'thinking' ? 'thinking' |
| 170 | : group === 'responding' ? 'assistant_message' : 'operation', categoryFor(kind), at, agentId); |
| 171 | this.active.set(key, { key, activityKind: kind, startedAtMs: at, event, group }); |
| 172 | } |
| 173 | |
| 174 | private pulse(key: string, at: number): void { |
| 175 | const span = this.active.get(key); |
| 176 | if (!span) return; |
| 177 | // A resumed stream does not assert tape coverage across its silent |
| 178 | // interval; the span itself (and its start) stays active. |
| 179 | if (at - span.event.endTime > PET_BIN_MS * 2) { |
| 180 | const event = this.add(span.event.name, span.event.category, at, span.event.agentId, true); |
| 181 | this.active.set(key, { ...span, event }); |
| 182 | } else { |
| 183 | span.event.endTime = at; |
| 184 | } |
| 185 | } |
| 186 | |
| 187 | private finish(key: string, at: number, status: 'success' | 'unknown' = 'success'): ActiveSpan | undefined { |
| 188 | if (!this.active.has(key)) return undefined; |
| 189 | this.pulse(key, at); |
| 190 | const span = this.active.get(key); |
| 191 | if (!span) return undefined; |
| 192 | span.event.endTime = at; |
| 193 | span.event.status = status; |
| 194 | this.active.delete(key); |
| 195 | return span; |
| 196 | } |
| 197 | |
| 198 | private addOnce(set: Set<string>, id: string, maximum: number): boolean { |
| 199 | if (set.has(id)) return false; |
| 200 | set.add(id); |
| 201 | while (set.size > maximum) set.delete(set.values().next().value as string); |
| 202 | return true; |
| 203 | } |
| 204 | |
| 205 | /** Transactional batch copy; validation errors cannot accept half a batch. */ |
| 206 | clone(): PetEngineTelemetry { |
| 207 | const next = new PetEngineTelemetry(); |
| 208 | const copy = <T>(value: T): T => JSON.parse(JSON.stringify(value)) as T; |
| 209 | next.events = copy(this.events); |
| 210 | const spans = new Map(next.events.map(event => [event.id, event])); |
| 211 | next.active = new Map(Array.from(this.active, ([key, span]) => [key, { |
| 212 | ...span, |
| 213 | event: spans.get(span.event.id) ?? copy(span.event), |
| 214 | }])); |
| 215 | next.waiting = this.waiting ? spans.get(this.waiting.id) ?? copy(this.waiting) : undefined; |
| 216 | next.sequence = this.sequence; |
| 217 | next.lastTime = this.lastTime; |
| 218 | next.lastObservedAt = this.lastObservedAt; |
| 219 | next.turnId = this.turnId; |
| 220 | next.turnOutcome = this.turnOutcome; |
| 221 | next.terminalAt = this.terminalAt; |
| 222 | next.lastFailedTool = this.lastFailedTool ? { ...this.lastFailedTool } : undefined; |
| 223 | next.completedSpans = new Set(this.completedSpans); |
| 224 | next.completedTurns = new Set(this.completedTurns); |
| 225 | return next; |
| 226 | } |
| 227 | |
| 228 | /** Mutates in place: clock, freshness and the event window change before |
| 229 | * type-specific fields are checked. Callers that continue after a throw must |
| 230 | * observe on a clone() and swap it in on success, as PetNative does. */ |
| 231 | observe(value: unknown, at: number): void { |
| 232 | if (!Number.isFinite(at) || at < this.lastTime || at > PET_MAX_SECONDS * 1000) |
| 233 | throw new Error('Invalid Engine pet clock.'); |
| 234 | if (!value || typeof value !== 'object' || Array.isArray(value)) |
| 235 | throw new Error('Invalid Engine pet metadata.'); |
| 236 | const event = value as Record<string, unknown>; |
| 237 | const allowed = ['event', 'index', 'channel', 'span_id', 'activity_kind', 'outcome', 'id', 'worker_status', 'turn_id', 'turn_outcome']; |
| 238 | const stringFields = ['event', 'span_id', 'id', 'worker_status', 'turn_id', 'turn_outcome']; |
| 239 | if (Object.keys(event).some(key => !allowed.includes(key)) || typeof event.event !== 'string' |
| 240 | || Object.values(event).some(v => typeof v === 'string' && v.length > 256) |
| 241 | || stringFields.some(key => event[key] !== undefined && typeof event[key] !== 'string') |
| 242 | || event.channel !== undefined && !['text', 'reasoning'].includes(event.channel as string) |
| 243 | || event.activity_kind !== undefined && !ACTIVITY_KINDS.includes(event.activity_kind as OwnerActivityKind) |
| 244 | || event.outcome !== undefined && (typeof event.outcome !== 'string' |
| 245 | || (event.event === 'approval_resolved' |
| 246 | ? !APPROVAL_OUTCOMES.includes(event.outcome as typeof APPROVAL_OUTCOMES[number]) |
| 247 | : event.event === 'operation_activity_completed' |
| 248 | ? !OPERATION_OUTCOMES.includes(event.outcome as OwnerOperationOutcome) |
| 249 | : true)) |
| 250 | || event.turn_outcome !== undefined && !TURN_OUTCOMES.includes(event.turn_outcome as OwnerTurnOutcome) |
| 251 | || event.index !== undefined && (!Number.isSafeInteger(event.index) || (event.index as number) < 0)) |
| 252 | throw new Error('Invalid Engine pet metadata fields.'); |
| 253 | |
| 254 | this.lastTime = at; |
| 255 | this.lastObservedAt = at; |
| 256 | this.events = this.events.filter(item => item.endTime >= at - 12_800); |
| 257 | const required = (key: string): string => { |
| 258 | const value = event[key]; |
| 259 | if (typeof value !== 'string' || !value) throw new Error(`Missing Engine ${key}.`); |
| 260 | return value; |
| 261 | }; |
| 262 | const index = (): string => { |
| 263 | if (!Number.isSafeInteger(event.index)) throw new Error('Missing Engine index.'); |
| 264 | return String(event.index); |
| 265 | }; |
| 266 | const activityKind = (): OwnerActivityKind => { |
| 267 | if (!ACTIVITY_KINDS.includes(event.activity_kind as OwnerActivityKind)) |
| 268 | throw new Error('Missing Engine activity kind.'); |
| 269 | return event.activity_kind as OwnerActivityKind; |
| 270 | }; |
| 271 | const startMessage = (key: string, kind: OwnerActivityKind, group: SpanGroup): void => { |
| 272 | this.start(key, kind, group, at); |
| 273 | this.waiting = undefined; |
| 274 | }; |
| 275 | |
| 276 | switch (event.event) { |
| 277 | case 'turn_started': { |
| 278 | const id = required('turn_id'); |
| 279 | if (this.turnId !== id) { |
| 280 | this.active.clear(); |
| 281 | this.waiting = undefined; |
| 282 | this.lastFailedTool = undefined; |
| 283 | this.turnId = id; |
| 284 | this.turnOutcome = undefined; |
| 285 | this.terminalAt = undefined; |
| 286 | } |
| 287 | break; |
| 288 | } |
| 289 | case 'message_started': startMessage(`message:${index()}`, 'responding', 'responding'); break; |
| 290 | case 'thinking_started': startMessage(`thinking:${index()}`, 'thinking', 'thinking'); break; |
| 291 | case 'response_delta': { |
| 292 | const reasoning = event.channel === 'reasoning'; |
| 293 | const key = `${reasoning ? 'thinking' : 'message'}:${index()}`; |
| 294 | if (!this.active.has(key)) this.start(key, reasoning ? 'thinking' : 'responding', reasoning ? 'thinking' : 'responding', at); |
| 295 | else this.pulse(key, at); |
| 296 | this.waiting = undefined; |
| 297 | break; |
| 298 | } |
| 299 | case 'message_complete': this.finish(`message:${index()}`, at); break; |
| 300 | case 'thinking_complete': this.finish(`thinking:${index()}`, at); break; |
| 301 | case 'operation_activity_started': { |
| 302 | const spanId = required('span_id'); |
| 303 | const kind = activityKind(); |
| 304 | if (!this.completedSpans.has(spanId)) this.start(`operation:${spanId}`, kind, 'operation', at); |
| 305 | this.waiting = undefined; |
| 306 | break; |
| 307 | } |
| 308 | case 'operation_activity_completed': { |
| 309 | const spanId = required('span_id'); |
| 310 | const kind = activityKind(); |
| 311 | if (!OPERATION_OUTCOMES.includes(event.outcome as OwnerOperationOutcome)) |
| 312 | throw new Error('Missing Engine operation outcome.'); |
| 313 | const outcome = event.outcome as OwnerOperationOutcome; |
| 314 | // A failure is recorded once, as its own onset event below; marking |
| 315 | // the span `error` too would count it twice on the tape. |
| 316 | const completed = this.finish(`operation:${spanId}`, at, |
| 317 | outcome === 'succeeded' ? 'success' : 'unknown'); |
| 318 | if (completed && this.addOnce(this.completedSpans, spanId, 4096)) { |
| 319 | if (outcome === 'failed') { |
| 320 | this.add('operation_failed', 'error', at).status = 'error'; |
| 321 | this.lastFailedTool = { activityKind: kind, failedAtMs: at }; |
| 322 | } |
| 323 | if (outcome === 'denied' || outcome === 'cancelled') this.waiting = undefined; |
| 324 | } |
| 325 | break; |
| 326 | } |
| 327 | case 'tool_call_heartbeat': |
| 328 | for (const [key, span] of this.active) if (span.group === 'operation') this.pulse(key, at); |
| 329 | break; |
| 330 | case 'approval_resolved': { |
| 331 | required('id'); |
| 332 | if (!APPROVAL_OUTCOMES.includes(event.outcome as typeof APPROVAL_OUTCOMES[number])) |
| 333 | throw new Error('Invalid Engine approval outcome.'); |
| 334 | this.waiting = undefined; |
| 335 | break; |
| 336 | } |
| 337 | case 'agent_spawned': { |
| 338 | const id = required('id'); |
| 339 | this.start(`agent:${id}`, 'delegating', 'agent', at, id); |
| 340 | break; |
| 341 | } |
| 342 | case 'agent_progress': { |
| 343 | const id = required('id'); |
| 344 | const key = `agent:${id}`; |
| 345 | if (['completed', 'failed', 'cancelled', 'interrupted', 'budget_exhausted'].includes(event.worker_status as string)) { |
| 346 | this.finish(key, at); |
| 347 | break; |
| 348 | } |
| 349 | if (!this.active.has(key)) this.start(key, 'delegating', 'agent', at, id); |
| 350 | else this.pulse(key, at); |
| 351 | break; |
| 352 | } |
| 353 | case 'agent_complete': this.finish(`agent:${required('id')}`, at); break; |
| 354 | case 'approval_required': case 'user_input_required': { |
| 355 | required('id'); |
| 356 | if (!this.waiting) { |
| 357 | this.waiting = this.add('human_request', 'human', at); |
| 358 | this.waiting.status = 'pending'; |
| 359 | } |
| 360 | break; |
| 361 | } |
| 362 | case 'turn_complete': { |
| 363 | const outcome = event.turn_outcome; |
| 364 | if (!TURN_OUTCOMES.includes(outcome as OwnerTurnOutcome)) |
| 365 | throw new Error('Missing Engine turn outcome.'); |
| 366 | const id = typeof event.turn_id === 'string' && event.turn_id.length ? event.turn_id : undefined; |
| 367 | this.active.clear(); |
| 368 | this.waiting = undefined; |
| 369 | this.turnId = id; |
| 370 | // An outcome belongs to a turn. `/purge`, an edit rejection or a |
| 371 | // session switch mid-turn completes with no turn id; recording the |
| 372 | // outcome alone would break the projection invariant the Rust |
| 373 | // contract checks (`turn_outcome` requires `turn_id`). |
| 374 | this.turnOutcome = id ? outcome as OwnerTurnOutcome : undefined; |
| 375 | this.terminalAt = at; |
| 376 | if (outcome === 'completed' && id && this.addOnce(this.completedTurns, id, 256)) |
| 377 | this.add('turn_completed', 'communication', at).status = 'success'; |
| 378 | break; |
| 379 | } |
| 380 | default: throw new Error('Unsupported Engine pet event.'); |
| 381 | } |
| 382 | } |
| 383 | |
| 384 | /** The typed shell may extend a request already witnessed in Engine events, |
| 385 | * but a mid-turn attach cannot invent a NeedsYou state. */ |
| 386 | confirmWaiting(at: number, waiting: boolean): void { |
| 387 | if (!waiting) { |
| 388 | this.waiting = undefined; |
| 389 | return; |
| 390 | } |
| 391 | if (this.waiting && at >= this.waiting.endTime) { |
| 392 | this.waiting.endTime = at; |
| 393 | this.lastObservedAt = at; |
| 394 | this.lastTime = Math.max(this.lastTime, at); |
| 395 | } |
| 396 | } |
| 397 | |
| 398 | bucket(sequence: number): PetBucket { |
| 399 | const end = (sequence + 1) * PET_BIN_MS; |
| 400 | const input = this.events.filter(event => event.startTime < end && event.endTime >= end - 12_400) |
| 401 | .map(event => ({ ...event, endTime: Math.min(event.endTime, end) })); |
| 402 | return compilePetTelemetry(input, end, sequence)[0]; |
| 403 | } |
| 404 | } |
| 405 |