返回 CodeWhale
index.js
根目录 / npm / runtime-sdk / index.js
1 const DEFAULT_BASE_URL = "http://127.0.0.1:7878";
2
3 export class RuntimeApiError extends Error {
4 constructor(message, options = {}) {
5 super(message);
6 this.name = "RuntimeApiError";
7 this.status = options.status;
8 this.method = options.method;
9 this.path = options.path;
10 this.body = options.body;
11 }
12 }
13
14 export class RuntimeCapabilityError extends RuntimeApiError {
15 constructor(capability, message, options = {}) {
16 super(message, options);
17 this.name = "RuntimeCapabilityError";
18 this.capability = capability;
19 }
20 }
21
22 export class CodeWhaleRuntimeClient {
23 constructor(options = {}) {
24 this.baseUrl = normalizeBaseUrl(options.baseUrl ?? DEFAULT_BASE_URL);
25 this.token = options.token ?? null;
26 this.fetchImpl = options.fetch ?? globalThis.fetch;
27 if (typeof this.fetchImpl !== "function") {
28 throw new TypeError("CodeWhaleRuntimeClient requires a fetch implementation");
29 }
30 }
31
32 async createFleetRun(spec) {
33 return this.#jsonRequest("/v1/fleet/runs", {
34 method: "POST",
35 body: spec,
36 capability: "fleet_run_create",
37 });
38 }
39
40 async startFleetRun(runId) {
41 return this.#jsonRequest(`/v1/fleet/runs/${segment(runId)}/start`, {
42 method: "POST",
43 capability: "fleet_run_start",
44 });
45 }
46
47 async replayFleetEvents(runId, options = {}) {
48 const path = fleetEventPath(
49 `/v1/fleet/runs/${segment(runId)}/events/replay`,
50 options,
51 );
52 return this.#jsonRequest(path, {
53 capability: "fleet_event_replay",
54 });
55 }
56
57 async listFleetRuns() {
58 return this.#jsonRequest("/v1/fleet/runs");
59 }
60
61 async getFleetRun(runId) {
62 return this.#jsonRequest(`/v1/fleet/runs/${segment(runId)}`);
63 }
64
65 async listFleetWorkers(runId) {
66 return this.#jsonRequest(`/v1/fleet/runs/${segment(runId)}/workers`);
67 }
68
69 async getFleetWorker(workerId) {
70 return this.#jsonRequest(`/v1/fleet/workers/${segment(workerId)}`);
71 }
72
73 async interruptWorker(workerId) {
74 return this.#jsonRequest(`/v1/fleet/workers/${segment(workerId)}/interrupt`, {
75 method: "POST",
76 });
77 }
78
79 async stopWorker(workerId) {
80 return this.#jsonRequest(`/v1/fleet/workers/${segment(workerId)}/stop`, {
81 method: "POST",
82 });
83 }
84
85 async restartWorker(workerId) {
86 return this.#jsonRequest(`/v1/fleet/workers/${segment(workerId)}/restart`, {
87 method: "POST",
88 });
89 }
90
91 async stopFleetRun(runId) {
92 return this.#jsonRequest(`/v1/fleet/runs/${segment(runId)}/stop`, {
93 method: "POST",
94 });
95 }
96
97 async *fleetEvents(runId, options = {}) {
98 const path = fleetEventPath(
99 options.path ?? `/v1/fleet/runs/${segment(runId)}/events`,
100 options,
101 );
102 const response = await this.#rawRequest(path, {
103 method: "GET",
104 capability: "fleet_event_stream",
105 accept: "text/event-stream",
106 });
107 const contentType = response.headers.get("content-type") ?? "";
108 if (contentType.includes("application/json")) {
109 const payload = await response.json();
110 const events = Array.isArray(payload) ? payload : (payload.events ?? []);
111 for (const event of events) {
112 yield event;
113 }
114 return;
115 }
116 if (!response.body) {
117 throw new RuntimeApiError("Runtime API event response did not include a readable body", {
118 method: "GET",
119 path,
120 });
121 }
122 for await (const event of parseEventStream(response.body)) {
123 yield event;
124 }
125 }
126
127 /** Read the existing durable thread journal. This never starts a turn. */
128 async *threadEvents(threadId, options = {}) {
129 const query = new URLSearchParams();
130 for (const [key, value] of [["since_seq", options.sinceSeq], ["replay_limit", options.replayLimit]]) {
131 if (value === undefined) continue;
132 if (!Number.isSafeInteger(value) || value < 0) throw new TypeError(`${key} must be a nonnegative safe integer`);
133 query.set(key, String(value));
134 }
135 if (options.includeProgress !== undefined && typeof options.includeProgress !== "boolean")
136 throw new TypeError("includeProgress must be a boolean");
137 if (options.includeProgress) query.set("progress", "true");
138 const path = `/v1/threads/${segment(threadId)}/events?${query}`;
139 const response = await this.#rawRequest(path, {
140 method: "GET", capability: "thread_event_stream", accept: "text/event-stream",
141 signal: options.signal, redirect: "error",
142 });
143 if (!response.body || !/^text\/event-stream(?:;|$)/i.test(response.headers.get("content-type") ?? "")) {
144 await response.body?.cancel();
145 throw new RuntimeApiError("Runtime thread response is not an event stream", { method: "GET", path });
146 }
147 if (options.includeProgress && response.headers.get("x-codewhale-event-progress") !== "1") {
148 await response.body.cancel();
149 throw new RuntimeCapabilityError("thread_event_progress", "Runtime does not support thread replay progress", { method: "GET", path, status: 501 });
150 }
151 yield* parseEventStream(response.body, { maxFrameChars: 2 * 1024 * 1024, requireBoundary: true });
152 }
153
154 async #jsonRequest(path, options = {}) {
155 const response = await this.#rawRequest(path, options);
156 if (response.status === 204) {
157 return null;
158 }
159 return response.json();
160 }
161
162 async #rawRequest(path, options = {}) {
163 const method = options.method ?? "GET";
164 const base = new URL(this.baseUrl);
165 const url = new URL(path, base);
166 if (!["http:", "https:"].includes(url.protocol) || url.origin !== base.origin || url.username || url.password) {
167 throw new TypeError("Runtime API requests must stay on the configured HTTP(S) origin without URL credentials");
168 }
169 const headers = new Headers(options.headers);
170 headers.set("accept", options.accept ?? "application/json");
171 if (this.token) {
172 headers.set("authorization", `Bearer ${this.token}`);
173 }
174 const init = { method, headers };
175 if (options.signal) init.signal = options.signal;
176 // A redirect can change the request's destination or repeat a mutation.
177 // Authenticated requests must not leave the origin checked above.
178 if (headers.has("authorization")) init.redirect = "error";
179 else if (options.redirect) init.redirect = options.redirect;
180 if (options.body !== undefined) {
181 headers.set("content-type", "application/json");
182 init.body = JSON.stringify(options.body);
183 }
184
185 const response = await this.fetchImpl(url, init);
186 if (response.ok) {
187 return response;
188 }
189
190 const body = await readErrorBody(response);
191 const errorOptions = { status: response.status, method, path, body };
192 if (options.capability && [404, 405, 501].includes(response.status)) {
193 throw new RuntimeCapabilityError(
194 options.capability,
195 `Runtime API capability '${options.capability}' is not available at ${method} ${path}`,
196 errorOptions,
197 );
198 }
199 throw new RuntimeApiError(
200 `Runtime API request failed (${response.status}) for ${method} ${path}`,
201 errorOptions,
202 );
203 }
204 }
205
206 export function createRuntimeClient(options = {}) {
207 return new CodeWhaleRuntimeClient(options);
208 }
209
210 /** Distinguish a stream-end transport frame from thread journal/progress events. */
211 export function isThreadStreamEnd(event) {
212 return event.event === "stream.end";
213 }
214
215 function normalizeBaseUrl(value) {
216 return value.endsWith("/") ? value : `${value}/`;
217 }
218
219 function segment(value) {
220 if (value === null || value === undefined || String(value).trim() === "") {
221 throw new TypeError("Runtime API path segment must be a non-empty value");
222 }
223 return encodeURIComponent(String(value));
224 }
225
226 function fleetEventPath(path, options) {
227 const query = new URLSearchParams();
228 if (options.after !== undefined && options.after !== null && String(options.after) !== "") {
229 query.set("after", String(options.after));
230 }
231 if (options.limit !== undefined && options.limit !== null) {
232 query.set("limit", String(options.limit));
233 }
234 const encoded = query.toString();
235 if (!encoded) {
236 return path;
237 }
238 return `${path}${path.includes("?") ? "&" : "?"}${encoded}`;
239 }
240
241 async function readErrorBody(response) {
242 try {
243 const text = await response.text();
244 return text.length > 4096 ? `${text.slice(0, 4096)}...` : text;
245 } catch {
246 return "";
247 }
248 }
249
250 async function* parseEventStream(body, { maxFrameChars = Infinity, requireBoundary = false } = {}) {
251 const decoder = new TextDecoder("utf-8", { fatal: requireBoundary });
252 let buffer = "";
253 for await (const chunk of body) {
254 buffer += decoder.decode(chunk, { stream: true });
255 let boundary;
256 while ((boundary = eventStreamBoundary(buffer)) !== null) {
257 if (boundary.index > maxFrameChars) throw new Error("Runtime event frame exceeds the size limit");
258 const frame = buffer.slice(0, boundary.index);
259 buffer = buffer.slice(boundary.index + boundary.length);
260 const event = parseSseFrame(frame);
261 if (event !== undefined) {
262 yield event;
263 }
264 }
265 if (buffer.length > maxFrameChars) throw new Error("Runtime event frame exceeds the size limit");
266 }
267 buffer += decoder.decode();
268 if (requireBoundary && buffer.trim()) throw new Error("Runtime event stream ended inside a frame");
269 const event = parseSseFrame(buffer);
270 if (event !== undefined) {
271 yield event;
272 }
273 }
274
275 function eventStreamBoundary(buffer) {
276 const lf = buffer.indexOf("\n\n");
277 const crlf = buffer.indexOf("\r\n\r\n");
278 if (lf < 0 && crlf < 0) {
279 return null;
280 }
281 if (crlf >= 0 && (lf < 0 || crlf < lf)) {
282 return { index: crlf, length: 4 };
283 }
284 return { index: lf, length: 2 };
285 }
286
287 function parseSseFrame(frame) {
288 const lines = frame.split(/\r?\n/);
289 const eventName = lines
290 .find((line) => line.startsWith("event:"))
291 ?.slice("event:".length)
292 .trimStart();
293 const eventId = lines
294 .find((line) => line.startsWith("id:"))
295 ?.slice("id:".length)
296 .trimStart();
297 const data = lines
298 .filter((line) => line.startsWith("data:"))
299 .map((line) => line.slice("data:".length).trimStart())
300 .join("\n");
301 if (!data || data === "[DONE]") {
302 return undefined;
303 }
304 const parsed = JSON.parse(data);
305 if (parsed && typeof parsed === "object" && !Array.isArray(parsed)) {
306 if (eventName && parsed.event === undefined) {
307 parsed.event = eventName;
308 }
309 if (eventId && parsed.cursor === undefined) {
310 parsed.cursor = eventId;
311 }
312 }
313 return parsed;
314 }
315
315 lines JAVASCRIPT