返回 CodeWhale
pet-runtime.mjs
根目录 / pet / scripts / lib / pet-runtime.mjs
1 import { setTimeout as delay } from 'node:timers/promises';
2 import { privacyEvent, redact } from '../../dist/core/ingest.js';
3 import { CodewhaleRuntimeTrace, isCodewhaleRuntimeRecord, observeRuntimeRequests } from '../../dist/core/codewhale.js';
4
5 /** What each Runtime `stream.end` reason means, for the recorder's report. */
6 const STREAM_END_REASONS = Object.freeze({
7 replay_failed: 'Runtime could not read the thread history',
8 catch_up_failed: 'Runtime fell behind and could not catch up',
9 runtime_shutdown: 'Runtime is shutting down',
10 });
11
12 /** A read-only transport for the existing Runtime journal. All event meaning
13 * remains in Whalesong's importer and canonical pet bucketer. No raw journal,
14 * prompt, tool argument or bearer token is written into the pet recording. */
15 export async function followRuntime({ baseUrl, threadId, token, report = () => {} }) {
16 const url = new URL(baseUrl);
17 if (url.protocol !== 'http:' || !['127.0.0.1', '[::1]'].includes(url.hostname)
18 || url.username || url.password || url.pathname !== '/' || url.search || url.hash)
19 throw new Error('Pet Runtime input requires a plain HTTP loopback IP origin, without credentials or a path.');
20 if (typeof threadId !== 'string' || !threadId.trim() || threadId.length > 512)
21 throw new Error('Choose one Runtime --thread ID.');
22 let sdk;
23 try { sdk = await import('@codewhale/runtime-sdk'); }
24 catch { sdk = await import('../../../npm/runtime-sdk/index.js'); }
25 if (typeof sdk.CodeWhaleRuntimeClient.prototype.threadEvents !== 'function')
26 throw new Error('The local Runtime SDK needs threadEvents support.');
27 const client = new sdk.CodeWhaleRuntimeClient({ baseUrl: url.href, token });
28 const shutdown = new AbortController();
29 const trace = new CodewhaleRuntimeTrace('Codewhale Runtime', 250_000,
30 event => privacyEvent(event, 'metadata'), 64 * 1024 * 1024);
31 let cursor = 0, revision = 0, connected = false, fatal = false;
32 const done = (async () => {
33 let backoff = 250;
34 while (!shutdown.signal.aborted && !fatal) {
35 // Fifteen-second server heartbeats make a silent, half-open connection
36 // distinguishable from an idle journal. The timeout is driver time only.
37 const attempt = new AbortController();
38 const signal = AbortSignal.any([shutdown.signal, attempt.signal]);
39 let idleTimer;
40 const refresh = () => { clearTimeout(idleTimer); idleTimer = setTimeout(() => attempt.abort(), 45_000); };
41 refresh();
42 const fetchImpl = client.fetchImpl;
43 client.fetchImpl = async (input, init) => {
44 const response = await fetchImpl(input, init);
45 if (!response.body || !response.ok) return response;
46 // Also reject an older installed SDK that silently omits the requested
47 // progress option; it must not turn historical packets into live state.
48 if (response.headers.get('x-codewhale-event-progress') !== '1') {
49 await response.body.cancel();
50 const error = new Error('Runtime replay progress is unavailable.'); error.status = 501; throw error;
51 }
52 // Cancel the wrapped pipeline too: the original Response can be collected
53 // while its idle body is still being read through the replacement below.
54 const body = response.body.pipeThrough(new TransformStream({ transform(chunk, controller) { refresh(); controller.enqueue(chunk); } }), { signal });
55 return new Response(body, { status: response.status, headers: response.headers });
56 };
57 try {
58 for await (const record of client.threadEvents(threadId, { sinceSeq: cursor, signal, includeProgress: true })) {
59 if (record?.event === 'stream.progress') {
60 if (record.thread_id !== threadId || record.seq !== cursor || !['live', 'replaying'].includes(record.state))
61 throw new Error('Invalid Runtime replay progress.');
62 connected = record.state === 'live';
63 continue;
64 }
65 if (record?.event === 'stream.end') {
66 // Runtime ended the stream on purpose at our cursor. Resume from it,
67 // and say what Runtime said rather than blame the connection.
68 if (record.thread_id !== threadId || record.last_seq !== cursor || typeof record.retryable !== 'boolean')
69 throw new Error('Invalid Runtime stream end.');
70 const why = STREAM_END_REASONS[record.reason] ?? 'Runtime ended the stream';
71 if (!record.retryable) fatal = true;
72 report(fatal
73 ? `Runtime input stopped: ${why}, and the stream cannot resume. Recording remains unobserved.`
74 : `Runtime input paused: ${why}; reconnecting from the last cursor.`);
75 break;
76 }
77 if (!isCodewhaleRuntimeRecord(record) || record.thread_id !== threadId || !Number.isSafeInteger(record.seq) || record.seq < 0)
78 throw new Error('Invalid Runtime envelope.');
79 if (record.seq <= cursor) continue;
80 // Sequence numbers belong to Runtime, and need not be consecutive.
81 // Its predecessor cursor detects loss without inventing a new counter.
82 if (record.previous_seq !== undefined && record.previous_seq !== cursor)
83 throw new Error('Runtime predecessor cursor does not match.');
84 if (record.event !== 'item.delta') {
85 // The existing importer retains unfinished lifetimes and a recent
86 // recurrence window, not a second copy of the entire raw journal.
87 if (revision % 256 === 0) trace.prune(Date.now() - 16_000);
88 try { trace.append([redact(record)]); }
89 catch (error) { fatal = true; throw error; }
90 revision++;
91 }
92 cursor = record.seq; backoff = 250;
93 }
94 } catch (error) {
95 if ([400, 401, 403, 404, 405, 501].includes(error.status)) fatal = true;
96 if (!shutdown.signal.aborted) report(fatal
97 ? 'Runtime input stopped: check the thread, authentication, SDK/Runtime replay-progress support, or retained input limit. Recording remains unobserved.'
98 : 'Runtime input interrupted; recording unobserved gaps while reconnecting from the last cursor.');
99 } finally {
100 connected = false; clearTimeout(idleTimer); client.fetchImpl = fetchImpl;
101 }
102 if (!shutdown.signal.aborted && !fatal) {
103 await delay(backoff, undefined, { signal: shutdown.signal }).catch(() => {});
104 backoff = Math.min(8000, backoff * 2);
105 }
106 }
107 })();
108 return {
109 get connected() { return connected; },
110 get revision() { return revision; },
111 get cursor() { return cursor; },
112 get retainedEvents() { return trace.retainedEvents; },
113 get retainedBytes() { return trace.retainedBytes; },
114 snapshot(observedThrough = Date.now()) {
115 if (!connected) return undefined;
116 trace.prune(observedThrough - 16_000);
117 if (!trace.retainedEvents) return undefined;
118 return observeRuntimeRequests({ ...trace.snapshot(), privacy: 'metadata' }, observedThrough);
119 },
120 async close() { shutdown.abort(); await done; },
121 };
122 }
123
123 lines Plain Text