返回 CodeWhale
codewhale.ts
根目录 / pet / src / core / codewhale.ts
1 /** Read-only Codewhale session/runtime adapter. Does not record, mutate, or own receipts. */
2 import { CATEGORIES, errorOnsetOf, type Category, type Status, type Trace, type WhaleEvent } from './model.js';
3
4 const PAYLOAD_LIMIT = 2000;
5 const COLLAPSED_SPAN_MS = 1000;
6 const ENVELOPE_MIN_MS = 60_000;
7
8 type Obj = Record<string, any>;
9 const obj = (v: unknown): Obj => v !== null && typeof v === 'object' && !Array.isArray(v) ? v as Obj : {};
10 const str = (v: unknown): string | undefined => typeof v === 'string' && v.length ? v : undefined;
11 const num = (v: unknown): number | undefined => typeof v === 'number' && Number.isFinite(v) ? v : undefined;
12
13 export function isCodewhaleSession(value: unknown): boolean {
14 const root = obj(value);
15 const metadata = obj(root.metadata);
16 if (!str(metadata.id)) return false;
17 if (root.format === 'whalesong.evidence/v1' || Array.isArray(root.resourceSpans) || root.schemaVersion === 1) return false;
18 const journal = obj(root.journal);
19 return Array.isArray(root.messages) || Array.isArray(journal.entries);
20 }
21
22 export function isCodewhaleRuntimeRecord(value: unknown): boolean {
23 const rec = obj(value);
24 return Number.isSafeInteger(rec.seq) && rec.seq >= 0 && typeof rec.event === 'string' && !!rec.event
25 && typeof rec.thread_id === 'string' && !!rec.thread_id && rec.timestamp != null;
26 }
27
28 export function isCodewhaleRuntimeDocument(value: unknown): boolean {
29 if (!Array.isArray(value) || !value.length) return false;
30 const n = Math.min(value.length, 8);
31 let hits = 0;
32 for (let i = 0; i < n; i++) if (isCodewhaleRuntimeRecord(value[i])) hits++;
33 return hits === n;
34 }
35
36 function clip(value: unknown): unknown {
37 if (value == null) return value;
38 const text = typeof value === 'string' ? value : JSON.stringify(value);
39 if (text.length <= PAYLOAD_LIMIT) return typeof value === 'string' ? value : JSON.parse(text);
40 return `${text.slice(0, PAYLOAD_LIMIT)}…[truncated ${text.length - PAYLOAD_LIMIT} source bytes]`;
41 }
42
43 function parseTime(value: unknown): number | undefined {
44 if (typeof value === 'number' && Number.isFinite(value)) return value;
45 if (typeof value !== 'string' || !value) return undefined;
46 const ms = Date.parse(value);
47 return Number.isFinite(ms) ? ms : undefined;
48 }
49
50 function statusOf(value: unknown, isError?: boolean): Status {
51 if (isError === true) return 'error';
52 if (isError === false) return 'success';
53 const s = String(value ?? '').toLowerCase();
54 if (s === 'completed' || s === 'success' || s === 'ok') return 'success';
55 if (s === 'failed' || s === 'error' || s === 'errored') return 'error';
56 if (s === 'canceled' || s === 'cancelled' || s === 'interrupted') return 'error';
57 if (s === 'in_progress' || s === 'running') return 'running';
58 if (s === 'pending') return 'pending';
59 return 'unknown';
60 }
61
62 function classify(name: string): Category {
63 const n = name.toLowerCase();
64 if (/exception|^error\b/.test(n)) return 'error';
65 if (/spawn|fork|subagent|^agent$/.test(n)) return 'agent';
66 if (/message\.send|handoff|agent\.message|assistant_message/.test(n)) return 'communication';
67 if (/retrieve|retrieval|context|embedding|vector|memory|rag/.test(n)) return 'memory';
68 if (/browser|navigate|screenshot|click|playwright/.test(n)) return 'browser';
69 if (/read_file|write_file|list_dir|^read$|^write$|^edit$|glob|grep|file\.|filesystem/.test(n)) return 'filesystem';
70 if (/bash|exec|shell|run_test|cargo|pytest|compile/.test(n)) return 'code';
71 if (/reason|thinking|completion|generate|chat|llm/.test(n)) return 'reasoning';
72 if (/http|request|api|fetch|network|mcp_/.test(n)) return 'network';
73 if (/user_message|human|approval/.test(n)) return 'human';
74 if (/orchestrat|workflow|phase|join|session|thread|turn|todo|plan|operate_contract|status/.test(n)) return 'orchestration';
75 if (/tool/.test(n)) return 'tool';
76 return CATEGORIES.includes(n as Category) ? n as Category : 'other';
77 }
78 export function toolCategory(name: string): Category {
79 const category = classify(name);
80 return category === 'other' ? 'tool' : category;
81 }
82
83 function pointer(source: string, ids: Obj): Obj {
84 return { format: source, ...ids };
85 }
86
87 function titleOfSession(metadata: Obj, filename: string): string {
88 const title = str(metadata.title)?.replace(/<[^>]+>/g, ' ').replace(/\s+/g, ' ').trim();
89 if (title && !title.startsWith('codewhale:runtime_event')) return title.slice(0, 120);
90 return `Codewhale session · ${(str(metadata.id) ?? filename).slice(0, 8)}`;
91 }
92
93 function activeJournalEntries(journal: Obj): { entries: Obj[]; warnings: string[] } {
94 const entries = Array.isArray(journal.entries) ? journal.entries.map(obj) : [];
95 const leaf = str(journal.leaf_id);
96 if (!leaf || !entries.length) return { entries, warnings: [] };
97 const byId = new Map(entries.filter(e => str(e.id)).map(e => [e.id as string, e]));
98 const chain: Obj[] = [];
99 const seen = new Set<string>();
100 let id: string | undefined = leaf;
101 while (id && !seen.has(id)) {
102 seen.add(id);
103 const entry = byId.get(id);
104 if (!entry) break;
105 chain.push(entry);
106 id = str(entry.parent_id);
107 }
108 if (!chain.length) return { entries, warnings: ['Journal leaf_id did not resolve; using append order instead of the active branch.'] };
109 if (chain.length < entries.length) {
110 return {
111 entries: chain.reverse(),
112 warnings: [`Active journal branch has ${chain.length} of ${entries.length} entries. Forked history was not invented into the timeline.`],
113 };
114 }
115 return { entries: chain.reverse(), warnings: [] };
116 }
117
118 function collapsedTimestamps(entries: Obj[], created?: number, updated?: number): boolean {
119 const times = entries.map(e => parseTime(e.created_at)).filter((n): n is number => n !== undefined);
120 if (times.length < 2) return false;
121 const span = Math.max(...times) - Math.min(...times);
122 const envelope = created !== undefined && updated !== undefined ? updated - created : 0;
123 return envelope >= ENVELOPE_MIN_MS && span < COLLAPSED_SPAN_MS;
124 }
125
126 function pushEvent(events: WhaleEvent[], event: WhaleEvent): void {
127 events.push(event);
128 }
129
130 export function fromCodewhaleSession(document: unknown, filename = 'Codewhale session', maxEvents = 250_000): Trace {
131 const root = obj(document);
132 const metadata = obj(root.metadata);
133 const sessionId = str(metadata.id) ?? filename;
134 const journal = obj(root.journal);
135 const { entries, warnings } = activeJournalEntries(journal);
136 const sourceEntries: Obj[] = entries.length ? entries : (Array.isArray(root.messages) ? root.messages.map((message: unknown, i: number) => ({ id: `${sessionId}/message/${i}`, kind: 'message', message })) : []);
137 if (!sourceEntries.length) throw new Error('Codewhale session contains no journal entries or messages.');
138 const created = parseTime(metadata.created_at);
139 const updated = parseTime(metadata.updated_at);
140 const orderOnly = collapsedTimestamps(sourceEntries, created, updated);
141 if (orderOnly) {
142 warnings.push('Journal created_at values are collapsed to last-save time, not execution time. The time axis is journal order (1 ms per emitted event), not wall-clock duration. Gap, burst, and cycle-period findings are not execution-time claims.');
143 } else {
144 const times = sourceEntries.map(e => parseTime(e.created_at)).filter((n): n is number => n !== undefined);
145 if (!times.length) warnings.push('Journal entries have no usable timestamps. The time axis is journal order.');
146 }
147
148 const events: WhaleEvent[] = [];
149 const pending = new Map<string, number>();
150 // One journal entry may hold any number of blocks, so bound every emitted
151 // event, not every entry. Identities stay unique like the generic importer:
152 // a repeated tool call id would otherwise replace the earlier call.
153 const identities = new Set<string>();
154 const push = (event: WhaleEvent): void => {
155 if (events.length >= maxEvents) throw new Error(`Import exceeds the ${maxEvents.toLocaleString()} event limit.`);
156 if (identities.has(event.id)) throw new Error(`Duplicate event identity (${sessionId}, ${event.id}). Import cancelled.`);
157 identities.add(event.id); pushEvent(events, event);
158 };
159 let seq = 0;
160 const originWall = orderOnly ? undefined : sourceEntries.map(e => parseTime(e.created_at)).find((n): n is number => n !== undefined);
161 const agentId = 'parent';
162 const model = str(metadata.model);
163 const provider = str(metadata.model_provider);
164
165 const when = (entry: Obj, fallback: number): { start: number; open: boolean } => {
166 if (orderOnly || originWall === undefined) return { start: fallback, open: false };
167 const t = parseTime(entry.created_at);
168 if (t === undefined) return { start: fallback, open: true };
169 return { start: t - originWall, open: false };
170 };
171
172 for (const entry of sourceEntries) {
173 const entryId = str(entry.id) ?? `${sessionId}/entry/${seq}`;
174 const message = obj(entry.message ?? (entry.kind === 'message' ? entry : {}));
175 const role = str(message.role) ?? (str(entry.kind) === 'user' ? 'user' : str(entry.kind) === 'assistant' ? 'assistant' : undefined);
176 const blocks: Obj[] = Array.isArray(message.content) ? message.content.map(obj) : [];
177 if (!blocks.length) {
178 const text = str(entry.text) ?? str(message.text);
179 if (text) blocks.push({ type: role === 'user' ? 'text' : 'text', text });
180 }
181 if (!blocks.length) continue;
182 const parentEventId = events.length ? events[events.length - 1]!.id : undefined;
183 for (const block of blocks) {
184 const t = when(entry, seq);
185 const idBase = `${entryId}/${seq}`;
186 const type = str(block.type) ?? 'text';
187 const raw = pointer('codewhale.session/v1', { sessionId, entryId, seq, blockType: type, toolUseId: block.id ?? block.tool_use_id });
188 if (type === 'tool_use' || type === 'server_tool_use') {
189 const tool = str(block.name) ?? 'tool';
190 const callId = str(block.id) ?? idBase;
191 const started = tool === 'agent' && obj(block.input).action === 'start';
192 const event: WhaleEvent = {
193 schemaVersion: 1, id: callId, traceId: sessionId, parentId: parentEventId,
194 startTime: t.start, endTime: t.start, openEnded: true,
195 agentId, name: started ? 'agent.spawn' : tool, tool, category: toolCategory(tool),
196 subtype: started ? 'fork' : undefined, model, provider,
197 status: 'running', attributes: { 'codewhale.entry_id': entryId, 'codewhale.seq': seq, 'tool.name': tool },
198 payload: { arguments: clip(block.input) }, raw,
199 };
200 pending.set(callId, events.length);
201 push(event);
202 } else if (type === 'tool_result') {
203 const callId = str(block.tool_use_id);
204 const isError = block.is_error === true;
205 const target = callId !== undefined ? pending.get(callId) : undefined;
206 if (target !== undefined) {
207 const prior = events[target]!;
208 // Skewed or missing entry times never produce a negative duration.
209 if (t.start < prior.startTime) warnings.push(`Tool ${callId} result is timestamped before its call; its duration is recorded as zero.`);
210 prior.endTime = Math.max(prior.startTime, t.start);
211 prior.openEnded = false;
212 prior.status = statusOf('completed', isError);
213 prior.payload = { ...(obj(prior.payload)), result: clip(block.content) };
214 prior.attributes = { ...prior.attributes, 'codewhale.result_entry_id': entryId };
215 pending.delete(callId as string);
216 } else {
217 push({
218 schemaVersion: 1, id: idBase, traceId: sessionId, parentId: callId ?? parentEventId,
219 startTime: t.start, endTime: t.start, agentId,
220 name: 'tool_result', category: 'tool', model, provider,
221 status: statusOf(undefined, isError),
222 attributes: { 'codewhale.entry_id': entryId, 'codewhale.seq': seq, tool_use_id: callId },
223 payload: { result: clip(block.content) }, raw,
224 });
225 }
226 } else if (type === 'thinking') {
227 push({
228 schemaVersion: 1, id: idBase, traceId: sessionId, parentId: parentEventId,
229 startTime: t.start, endTime: t.start, agentId, name: 'thinking', category: 'reasoning',
230 model, provider, status: 'success',
231 attributes: { 'codewhale.entry_id': entryId, 'codewhale.seq': seq },
232 payload: { thinking: clip(block.thinking ?? block.text) }, raw,
233 });
234 } else {
235 const text = str(block.text) ?? '';
236 const operate = text.includes('codewhale:runtime_event');
237 const user = role === 'user' || role === 'User';
238 push({
239 schemaVersion: 1, id: idBase, traceId: sessionId, parentId: parentEventId,
240 startTime: t.start, endTime: t.start, agentId,
241 name: operate ? 'operate_contract' : user ? 'user_message' : 'assistant_message',
242 category: operate ? 'orchestration' : user ? 'human' : 'communication',
243 model, provider, status: 'success',
244 attributes: { 'codewhale.entry_id': entryId, 'codewhale.seq': seq, role: role ?? 'unknown' },
245 payload: { text: clip(text) }, raw,
246 });
247 }
248 seq += 1;
249 }
250 }
251
252 if (!events.length) throw new Error('Codewhale session produced no inspectable events.');
253 for (const event of events) {
254 if (event.openEnded && event.tool) warnings.push(`Tool ${event.id} has no matching tool_result in this snapshot; duration remains unknown.`);
255 }
256 const base = events.reduce((m, e) => Math.min(m, e.startTime), events[0]!.startTime);
257 for (const event of events) { event.startTime -= base; event.endTime -= base; }
258 const cost = obj(metadata.cost);
259 const sessionCost = num(cost.session_cost_usd);
260 const duration = Math.max(1, events.reduce((m, e) => Math.max(m, e.endTime, e.startTime), 0));
261 const uniqueWarnings = [...new Set(warnings)];
262 return {
263 id: sessionId,
264 name: titleOfSession(metadata, filename),
265 events,
266 duration,
267 originTime: orderOnly ? 'journal-order' : (str(metadata.created_at) ?? `${base} ms`),
268 source: 'codewhale',
269 privacy: 'redact',
270 warnings: uniqueWarnings,
271 metadata: {
272 sourceFormat: 'codewhale.session/v1',
273 timeBasis: orderOnly || originWall === undefined ? 'journal-order' : 'wall-clock',
274 sourceFilename: filename,
275 sessionId,
276 model,
277 provider,
278 workspace: metadata.workspace,
279 mode: metadata.mode,
280 envelopeCreatedAt: metadata.created_at,
281 envelopeUpdatedAt: metadata.updated_at,
282 cumulativeTurnSecs: metadata.cumulative_turn_secs,
283 messageCount: metadata.message_count,
284 journalEntries: sourceEntries.length,
285 totalTokens: metadata.total_tokens,
286 sessionCostUsd: sessionCost,
287 pricedTurns: cost.priced_turns,
288 unpricedTurns: cost.unpriced_turns,
289 runtimeStore: metadata.runtime_store,
290 timeUnit: 'ms',
291 },
292 };
293 }
294
295 function itemToolName(item: Obj, payload: Obj): string | undefined {
296 const named = str(payload.tool) ?? str(item.tool) ?? str(item.name);
297 if (named) return named;
298 if (str(item.kind) !== 'tool_call') return undefined;
299 const head = str(item.summary)?.split(':')[0]?.trim();
300 if (head && head.length < 80 && !/\s/.test(head)) return head;
301 return undefined;
302 }
303
304 function itemCategory(kind: string, tool?: string): Category {
305 if (kind === 'user_message') return 'human';
306 if (kind === 'agent_reasoning') return 'reasoning';
307 if (kind === 'agent_message') return 'communication';
308 if (kind === 'status') return 'orchestration';
309 if (kind === 'tool_call' && tool) return toolCategory(tool);
310 if (kind === 'tool_call') return 'tool';
311 return classify(kind);
312 }
313
314 /** Incremental form of the existing Runtime importer. File imports and live
315 * recording share this exact lifecycle parser; only a live driver retires old
316 * completed events after it has recorded their projection. */
317 export class CodewhaleRuntimeTrace {
318 private readonly events: WhaleEvent[] = [];
319 private readonly open = new Map<string, WhaleEvent>();
320 private readonly requests = new Map<string, WhaleEvent>();
321 private readonly sizes = new Map<WhaleEvent, number>();
322 private bytes = 0;
323 private recordCount = 0;
324 private skippedDeltas = 0;
325 private origin: number | undefined;
326 private model: string | undefined;
327 private threadId: string | undefined;
328 private threadName: string | undefined;
329 constructor(private readonly filename = 'Codewhale runtime', private readonly maxEvents = 250_000,
330 private readonly project: (event: WhaleEvent) => WhaleEvent = event => event,
331 private readonly maxBytes = Infinity) {}
332 get retainedEvents(): number { return this.events.length; }
333 get retainedBytes(): number { return this.bytes; }
334 private measure(event: WhaleEvent, proposed = event): void {
335 const safe = this.project(proposed);
336 if (this.maxBytes !== Infinity) {
337 const size = new TextEncoder().encode(JSON.stringify(safe)).length;
338 const total = this.bytes - (this.sizes.get(event) ?? 0) + size;
339 if (total > this.maxBytes) throw new Error('Runtime observation exceeds its retained input limit.');
340 this.bytes = total; this.sizes.set(event, size);
341 }
342 for (const key of Object.keys(event)) if (!Object.hasOwn(safe, key)) delete (event as unknown as Obj)[key];
343 Object.assign(event, safe);
344 }
345 private push(event: WhaleEvent): void {
346 if (this.events.length >= this.maxEvents) throw new Error(`Import exceeds the ${this.maxEvents.toLocaleString()} event limit.`);
347 this.measure(event); pushEvent(this.events, event);
348 }
349 /** Keep unfinished lifetimes plus the recent window needed by the bucketer's
350 * 12-second recurrence measure. A completion may still arrive for any open item. */
351 prune(beforeWall: number): void {
352 if (!Number.isFinite(beforeWall)) throw new Error('Invalid Runtime retention horizon.');
353 if (this.origin === undefined) return;
354 const cutoff = beforeWall - this.origin;
355 let keep = 0;
356 for (const event of this.events) {
357 if (event.openEnded || Math.max(event.endTime, errorOnsetOf(event)) >= cutoff) this.events[keep++] = event;
358 else {
359 this.bytes -= this.sizes.get(event) ?? 0; this.sizes.delete(event);
360 if (this.open.get(event.id) === event) this.open.delete(event.id);
361 }
362 }
363 this.events.length = keep;
364 }
365 append(records: unknown[]): void {
366 if (!records.length) return;
367 const { events, open, requests } = this;
368 const threadId = this.threadId ?? str(obj(records[0]).thread_id) ?? this.filename;
369 this.threadId = threadId;
370 let { origin, model, skippedDeltas } = this;
371 let threadName = this.threadName ?? threadId;
372 const stamp = (rec: Obj): number => {
373 const t = parseTime(rec.timestamp);
374 if (t === undefined) throw new Error(`Runtime event seq ${rec.seq} is missing a usable timestamp.`);
375 if (origin === undefined) origin = t;
376 return t - origin;
377 };
378 for (const raw of records) {
379 if (!isCodewhaleRuntimeRecord(raw)) throw new Error('Runtime import cancelled: a line is not a Codewhale runtime event record. No rows were skipped.');
380 this.recordCount++;
381 const rec = obj(raw);
382 if (rec.thread_id !== threadId) throw new Error('Runtime import contains multiple threads. Export one thread before importing.');
383 const eventName = rec.event as string;
384 if (eventName === 'item.delta') { skippedDeltas++; continue; }
385 const payload = obj(rec.payload);
386 const item = obj(payload.item);
387 const turn = obj(payload.turn);
388 const thread = obj(payload.thread);
389 const relative = stamp(rec);
390 const turnId = str(rec.turn_id) ?? str(payload.turn_id);
391 const itemId = str(rec.item_id) ?? str(item.id);
392 const agentId = 'parent';
393 if (str(thread.model)) model = str(thread.model);
394 if (str(turn.model)) model = str(turn.model) ?? model;
395
396 if (eventName === 'thread.started') {
397 model = str(thread.model) ?? model;
398 threadName = str(thread.id) ?? threadId;
399 this.push({
400 schemaVersion: 1, id: `thread:${threadId}`, traceId: threadId,
401 startTime: relative, endTime: relative, openEnded: true,
402 agentId, name: 'thread', category: 'orchestration', model, status: 'running',
403 attributes: { 'codewhale.seq': rec.seq, 'whalesong.container': true }, raw: rec,
404 });
405 continue;
406 }
407 if (eventName === 'turn.started' || eventName === 'turn.completed') {
408 const id = `turn:${turnId ?? rec.seq}`;
409 if (eventName === 'turn.completed') for (const [key, request] of requests) {
410 if (request.parentId !== id) continue;
411 this.measure(request, { ...request, endTime: Math.max(request.startTime, relative), openEnded: false, status: 'unknown' });
412 requests.delete(key);
413 }
414 const startWall = parseTime(turn.started_at) ?? parseTime(turn.created_at);
415 const endWall = parseTime(turn.ended_at);
416 const start = startWall !== undefined && origin !== undefined ? startWall - origin : relative;
417 const end = eventName === 'turn.completed' && endWall !== undefined && origin !== undefined ? endWall - origin : relative;
418 const usage = obj(turn.usage);
419 const existing = events.findIndex(e => e.id === id);
420 const next: WhaleEvent = {
421 schemaVersion: 1, id, traceId: threadId, parentId: `thread:${threadId}`,
422 startTime: start, endTime: Math.max(start, end), openEnded: eventName !== 'turn.completed',
423 agentId, name: 'turn', category: 'orchestration', model,
424 inputTokens: num(usage.input_tokens), outputTokens: num(usage.output_tokens),
425 status: statusOf(turn.status ?? payload.status), latency: num(turn.duration_ms),
426 attributes: { 'codewhale.seq': rec.seq, 'codewhale.turn_id': turnId, 'whalesong.container': true,
427 ...(statusOf(turn.status ?? payload.status) === 'error' ? { 'whalesong.error_onset_ms': relative } : {}) },
428 payload: { input_summary: clip(turn.input_summary) }, raw: rec,
429 };
430 if (existing >= 0) {
431 // The kept start can follow a skewed completion time; never end before it.
432 const prior = events[existing]!;
433 this.measure(prior, { ...next, startTime: prior.startTime, endTime: Math.max(prior.startTime, next.endTime) });
434 }
435 else this.push(next);
436 continue;
437 }
438 if (eventName === 'turn.lifecycle') continue;
439 if (['approval.required', 'approval.decided', 'approval.timeout', 'user_input.required', 'user_input.answered', 'user_input.canceled'].includes(eventName)) {
440 const kind = eventName.startsWith('approval.') ? 'approval' : 'user_input';
441 const requestId = str(payload[kind === 'approval' ? 'approval_id' : 'input_id']) ?? str(payload.id);
442 if (!requestId) throw new Error(`Runtime ${eventName} is missing its request identity.`);
443 const key = JSON.stringify([turnId ?? '', kind, requestId]);
444 const prior = requests.get(key), required = eventName.endsWith('.required');
445 if (required && prior) continue;
446 if (!required && prior) {
447 const next: WhaleEvent = { ...prior, attributes: { ...prior.attributes },
448 endTime: Math.max(prior.startTime, relative), openEnded: false,
449 status: eventName === 'approval.decided' || eventName === 'user_input.answered' ? 'success' : 'unknown' };
450 if (payload.auto === true) {
451 // Automatic consent has a receipt, but never asked the human to wait.
452 next.category = 'orchestration'; delete next.attributes['whalesong.waiting'];
453 next.attributes['whalesong.container'] = true;
454 }
455 this.measure(prior, next); requests.delete(key); continue;
456 }
457 const automatic = payload.auto === true;
458 const event: WhaleEvent = {
459 schemaVersion: 1, id: `request:${key}:${rec.seq}`, traceId: threadId,
460 parentId: turnId ? `turn:${turnId}` : undefined, startTime: relative, endTime: relative,
461 openEnded: required, agentId, name: eventName, category: automatic ? 'orchestration' : 'human',
462 status: required ? 'pending' : 'success', model,
463 attributes: { 'codewhale.seq': rec.seq, 'whalesong.waiting': required, 'whalesong.container': automatic }, raw: rec,
464 };
465 this.push(event);
466 if (required) requests.set(key, event);
467 continue;
468 }
469 if (eventName === 'tool_call.requested' || eventName === 'tool_call.canceled') {
470 const callId = str(payload.call_id) ?? `call:${rec.seq}`;
471 const tool = str(payload.tool);
472 const canceled = eventName === 'tool_call.canceled';
473 this.push({
474 schemaVersion: 1, id: `${eventName}:${callId}`, traceId: threadId, parentId: turnId ? `turn:${turnId}` : undefined,
475 startTime: relative, endTime: relative, agentId, name: tool ?? eventName, tool,
476 category: tool ? toolCategory(tool) : 'tool', model, status: canceled ? 'error' : 'pending',
477 attributes: { 'codewhale.seq': rec.seq, 'codewhale.turn_id': turnId, 'codewhale.call_id': callId, reason: payload.reason },
478 payload: { arguments: clip(payload.arguments) }, raw: rec,
479 });
480 continue;
481 }
482 if (eventName === 'item.started' || eventName === 'item.completed') {
483 const kind = str(item.kind) ?? 'item';
484 const tool = itemToolName(item, payload);
485 const id = itemId ?? `item:${rec.seq}`;
486 const startWall = parseTime(item.started_at);
487 const endWall = parseTime(item.ended_at);
488 const start = startWall !== undefined && origin !== undefined ? startWall - origin : relative;
489 const end = eventName === 'item.completed' && endWall !== undefined && origin !== undefined ? endWall - origin : relative;
490 const openEnded = eventName === 'item.started' && endWall === undefined;
491 const existing = open.get(id);
492 if (existing && eventName === 'item.completed') {
493 const prior = existing;
494 const status = statusOf(item.status);
495 this.measure(prior, { ...prior, endTime: Math.max(prior.startTime, end), openEnded: false, status,
496 attributes: { ...prior.attributes, ...(status === 'error' ? { 'whalesong.error_onset_ms': relative } : {}) },
497 payload: { summary: clip(item.summary), detail: clip(item.detail) } });
498 open.delete(id);
499 continue;
500 }
501 const event: WhaleEvent = {
502 schemaVersion: 1, id, traceId: threadId, parentId: turnId ? `turn:${turnId}` : undefined,
503 startTime: start, endTime: Math.max(start, end), openEnded,
504 agentId, name: tool ?? kind, tool, category: itemCategory(kind, tool), model,
505 status: statusOf(item.status ?? (eventName === 'item.started' ? 'running' : undefined)),
506 attributes: { 'codewhale.seq': rec.seq, 'codewhale.turn_id': turnId, 'codewhale.item_kind': kind,
507 ...(statusOf(item.status) === 'error' ? { 'whalesong.error_onset_ms': relative } : {}) },
508 payload: { summary: clip(item.summary), detail: clip(item.detail) }, raw: rec,
509 };
510 this.push(event);
511 if (eventName === 'item.started') open.set(id, event);
512 continue;
513 }
514 this.push({
515 schemaVersion: 1, id: `${eventName}:${rec.seq}`, traceId: threadId, parentId: turnId ? `turn:${turnId}` : undefined,
516 startTime: relative, endTime: relative, agentId, name: eventName, category: classify(eventName),
517 model, status: 'unknown', attributes: { 'codewhale.seq': rec.seq }, raw: rec,
518 });
519 }
520
521 this.origin = origin; this.model = model; this.threadName = threadName; this.skippedDeltas = skippedDeltas;
522 }
523 snapshot(): Trace {
524 const { events, open, requests, origin, model, skippedDeltas, filename } = this;
525 const threadId = this.threadId ?? filename, threadName = this.threadName ?? threadId;
526 const warnings: string[] = [];
527 if (skippedDeltas) warnings.push(`Dropped ${skippedDeltas.toLocaleString()} item.delta records; they are token stream fragments, not spans. Item start/end remain the source of duration.`);
528 for (const [id] of open) warnings.push(`Item ${id} started and never completed in this file; duration remains unknown.`);
529 for (const request of requests.values()) warnings.push(`Request ${request.id} has no terminal receipt; its duration remains unknown in this file.`);
530 if (!events.length) throw new Error('Codewhale runtime file contained only stream deltas or unreadable records.');
531 const base = events.reduce((m, e) => Math.min(m, e.startTime), events[0]!.startTime);
532 const normalized = events.map(event => ({ ...event, startTime: event.startTime - base, endTime: event.endTime - base,
533 attributes: { ...event.attributes, ...(event.attributes['whalesong.error_onset_ms'] !== undefined
534 ? { 'whalesong.error_onset_ms': errorOnsetOf(event) - base } : {}) } }));
535 return {
536 id: threadId,
537 name: `Codewhale runtime · ${threadName}`,
538 events: normalized,
539 duration: Math.max(1, normalized.reduce((m, e) => Math.max(m, e.endTime, e.startTime, e.status === 'error' ? errorOnsetOf(e) : 0), 0)),
540 originTime: origin !== undefined ? new Date(origin + base).toISOString() : '0 ms',
541 source: 'codewhale',
542 privacy: 'redact',
543 warnings: [...new Set(warnings)],
544 metadata: {
545 sourceFormat: 'codewhale.runtime-events/v2',
546 timeBasis: 'wall-clock',
547 sourceFilename: filename,
548 threadId,
549 model,
550 skippedDeltas,
551 recordCount: this.recordCount,
552 timeUnit: 'ms',
553 },
554 };
555 }
556 }
557
558 export function fromCodewhaleRuntime(records: unknown[], filename = 'Codewhale runtime', maxEvents = 250_000): Trace {
559 if (!records.length) throw new Error('Codewhale runtime event file is empty.');
560 const trace = new CodewhaleRuntimeTrace(filename, maxEvents);
561 trace.append(records); return trace.snapshot();
562 }
563
564 /** The journal owns request state until a matching terminal receipt. A live
565 * driver may confirm that state only while its cursor-checked stream is healthy.
566 * Ordinary open tool spans remain unknown-duration; no execution is inferred. */
567 export function observeRuntimeRequests(trace: Trace, observedThrough: number): Trace {
568 const origin = Date.parse(trace.originTime ?? '');
569 if (trace.metadata.sourceFormat !== 'codewhale.runtime-events/v2' || !Number.isFinite(origin)
570 || !Number.isFinite(observedThrough)) throw new Error('Invalid Runtime observation horizon.');
571 const at = observedThrough - origin;
572 const events = trace.events.map(e => e.openEnded && e.attributes['whalesong.waiting'] === true && at >= e.startTime
573 ? { ...e, endTime: at, openEnded: false } : e);
574 return { ...trace, events, duration: Math.max(trace.duration, at) };
575 }
576
576 lines TYPESCRIPT