392 lines
11 KiB
TypeScript
392 lines
11 KiB
TypeScript
import { KUBER_API_BASE_URL } from "../const";
|
|
import type { ApiProblemDetails } from "../shared/api";
|
|
import { KUBER_VERSION_HEADER } from "../shared/version";
|
|
import { readSession, type KuberSession } from "./session";
|
|
|
|
export type { ApiProblemDetails } from "../shared/api";
|
|
|
|
export const DEFAULT_API_TIMEOUT_MS = 30_000;
|
|
|
|
export type ApiRequestInit = RequestInit & {
|
|
/** JSON is serialized here so callers do not need to manage content headers. */
|
|
json?: unknown;
|
|
};
|
|
|
|
export type ApiRequestOptions = {
|
|
authenticated?: boolean;
|
|
baseUrl?: string;
|
|
session?: KuberSession;
|
|
/** Set to 0 to disable the deadline, primarily for long-lived streams. */
|
|
timeoutMs?: number;
|
|
/** Byte offset for resumable binary upload requests. */
|
|
uploadOffset?: number;
|
|
/** Overrides version observation for embedded clients and tests. */
|
|
onServerVersion?: (version: string | null) => void;
|
|
};
|
|
|
|
export type ApiUploadOptions = ApiRequestOptions & {
|
|
offset: number;
|
|
contentType?: string;
|
|
};
|
|
|
|
export class KuberApiError extends Error {
|
|
readonly code: string;
|
|
readonly requestId?: string;
|
|
readonly operationId?: string;
|
|
readonly problem: ApiProblemDetails;
|
|
|
|
constructor(message: string, status: number, problem?: ApiProblemDetails) {
|
|
super(message);
|
|
this.name = "KuberApiError";
|
|
this.status = status;
|
|
this.code = problem?.code ?? `HTTP_${status}`;
|
|
this.requestId = problem?.requestId;
|
|
this.operationId = problem?.operationId;
|
|
this.problem = {
|
|
...problem,
|
|
status,
|
|
title: problem?.title ?? message,
|
|
code: this.code,
|
|
};
|
|
}
|
|
|
|
readonly status: number;
|
|
}
|
|
|
|
type RequestDeadline = {
|
|
signal: AbortSignal;
|
|
clear: () => void;
|
|
};
|
|
|
|
function requestDeadline(
|
|
signal: AbortSignal | null | undefined,
|
|
timeoutMs = DEFAULT_API_TIMEOUT_MS,
|
|
): RequestDeadline {
|
|
if (!Number.isFinite(timeoutMs) || timeoutMs < 0) {
|
|
throw new RangeError("timeoutMs must be a finite non-negative number");
|
|
}
|
|
|
|
const controller = new AbortController();
|
|
const abortFromCaller = () => controller.abort(signal?.reason);
|
|
if (signal?.aborted) abortFromCaller();
|
|
else signal?.addEventListener("abort", abortFromCaller, { once: true });
|
|
|
|
const timer =
|
|
timeoutMs === 0
|
|
? undefined
|
|
: setTimeout(() => {
|
|
controller.abort(
|
|
new DOMException(
|
|
`API request timed out after ${timeoutMs}ms`,
|
|
"TimeoutError",
|
|
),
|
|
);
|
|
}, timeoutMs);
|
|
|
|
return {
|
|
signal: controller.signal,
|
|
clear: () => {
|
|
if (timer !== undefined) clearTimeout(timer);
|
|
signal?.removeEventListener("abort", abortFromCaller);
|
|
},
|
|
};
|
|
}
|
|
|
|
type FetchBody = NonNullable<RequestInit["body"]>;
|
|
|
|
function isBinaryBody(body: FetchBody): boolean {
|
|
return (
|
|
body instanceof ArrayBuffer ||
|
|
ArrayBuffer.isView(body) ||
|
|
(typeof Blob !== "undefined" && body instanceof Blob) ||
|
|
(typeof ReadableStream !== "undefined" && body instanceof ReadableStream)
|
|
);
|
|
}
|
|
|
|
async function requestHeaders(
|
|
init: ApiRequestInit,
|
|
options: ApiRequestOptions,
|
|
): Promise<Headers> {
|
|
const headers = new Headers(init.headers);
|
|
if (!headers.has("accept")) headers.set("accept", "application/json");
|
|
|
|
if (Object.hasOwn(init, "json")) {
|
|
headers.set("content-type", "application/json");
|
|
} else if (
|
|
init.body !== undefined &&
|
|
init.body !== null &&
|
|
!headers.has("content-type") &&
|
|
!isBinaryBody(init.body)
|
|
) {
|
|
// Preserve the original apiRequest convention: string bodies are JSON.
|
|
headers.set("content-type", "application/json");
|
|
}
|
|
|
|
if (options.uploadOffset !== undefined) {
|
|
if (
|
|
!Number.isSafeInteger(options.uploadOffset) ||
|
|
options.uploadOffset < 0
|
|
) {
|
|
throw new RangeError("uploadOffset must be a non-negative safe integer");
|
|
}
|
|
headers.set("upload-offset", String(options.uploadOffset));
|
|
}
|
|
|
|
if (options.authenticated !== false) {
|
|
const session = options.session ?? (await readSession());
|
|
if (!session) throw new Error("Not logged in. Run kuber login first.");
|
|
headers.set("authorization", `Bearer ${session.token}`);
|
|
}
|
|
return headers;
|
|
}
|
|
|
|
async function responseProblem(response: Response): Promise<ApiProblemDetails> {
|
|
let problem: ApiProblemDetails | undefined;
|
|
try {
|
|
const value: unknown = JSON.parse(await response.text());
|
|
if (value && typeof value === "object") {
|
|
problem = value as ApiProblemDetails;
|
|
}
|
|
} catch {
|
|
// The status and response headers still provide a stable error shape.
|
|
}
|
|
|
|
return {
|
|
...problem,
|
|
status: response.status,
|
|
title: problem?.title || response.statusText || `HTTP ${response.status}`,
|
|
code: problem?.code || `HTTP_${response.status}`,
|
|
requestId:
|
|
problem?.requestId ?? response.headers.get("x-request-id") ?? undefined,
|
|
operationId:
|
|
problem?.operationId ??
|
|
response.headers.get("x-operation-id") ??
|
|
undefined,
|
|
};
|
|
}
|
|
|
|
async function assertResponseOk(response: Response): Promise<void> {
|
|
if (response.ok) return;
|
|
const problem = await responseProblem(response);
|
|
throw new KuberApiError(
|
|
problem.detail || problem.message || problem.title,
|
|
response.status,
|
|
problem,
|
|
);
|
|
}
|
|
|
|
async function sendRequest(
|
|
path: string,
|
|
init: ApiRequestInit,
|
|
options: ApiRequestOptions,
|
|
signal: AbortSignal,
|
|
): Promise<Response> {
|
|
const headers = await requestHeaders(init, options);
|
|
const { json, ...requestInit } = init;
|
|
const body = Object.hasOwn(init, "json") ? JSON.stringify(json) : init.body;
|
|
const response = await fetch(
|
|
`${options.baseUrl ?? KUBER_API_BASE_URL}${path}`,
|
|
{
|
|
...requestInit,
|
|
body,
|
|
headers,
|
|
signal,
|
|
},
|
|
);
|
|
try {
|
|
options.onServerVersion?.(response.headers.get(KUBER_VERSION_HEADER));
|
|
} catch {
|
|
// Version discovery must never interrupt an API request.
|
|
}
|
|
return response;
|
|
}
|
|
|
|
export async function apiRequest<T>(
|
|
path: string,
|
|
init: ApiRequestInit = {},
|
|
options: ApiRequestOptions = {},
|
|
): Promise<T> {
|
|
const deadline = requestDeadline(init.signal, options.timeoutMs);
|
|
try {
|
|
const response = await sendRequest(path, init, options, deadline.signal);
|
|
await assertResponseOk(response);
|
|
|
|
if (
|
|
init.method?.toUpperCase() === "HEAD" ||
|
|
response.status === 204 ||
|
|
response.status === 205 ||
|
|
response.body === null
|
|
) {
|
|
return undefined as T;
|
|
}
|
|
|
|
const text = await response.text();
|
|
if (!text.trim()) return undefined as T;
|
|
return JSON.parse(text) as T;
|
|
} finally {
|
|
deadline.clear();
|
|
}
|
|
}
|
|
|
|
export function apiUpload<T = void>(
|
|
path: string,
|
|
body: Blob | ArrayBuffer | ArrayBufferView | ReadableStream<Uint8Array>,
|
|
options: ApiUploadOptions,
|
|
): Promise<T> {
|
|
const {
|
|
contentType = "application/octet-stream",
|
|
offset,
|
|
...requestOptions
|
|
} = options;
|
|
return apiRequest<T>(
|
|
path,
|
|
{
|
|
method: "PATCH",
|
|
headers: { "content-type": contentType },
|
|
body: body as FetchBody,
|
|
},
|
|
{ ...requestOptions, uploadOffset: offset },
|
|
);
|
|
}
|
|
|
|
/** Parses records as they arrive instead of buffering the complete response. */
|
|
export async function* apiStreamNdjson<T>(
|
|
path: string,
|
|
init: ApiRequestInit = {},
|
|
options: ApiRequestOptions = {},
|
|
): AsyncGenerator<T, void, void> {
|
|
const deadline = requestDeadline(init.signal, options.timeoutMs);
|
|
try {
|
|
const headers = new Headers(init.headers);
|
|
headers.set("accept", "application/x-ndjson");
|
|
const response = await sendRequest(
|
|
path,
|
|
{ ...init, headers },
|
|
options,
|
|
deadline.signal,
|
|
);
|
|
await assertResponseOk(response);
|
|
if (!response.body) return;
|
|
|
|
const reader = response.body.getReader();
|
|
try {
|
|
const decoder = new TextDecoder();
|
|
let buffer = "";
|
|
for (;;) {
|
|
const { done, value } = await reader.read();
|
|
buffer += decoder.decode(value, { stream: !done });
|
|
let start = 0;
|
|
let newline = buffer.indexOf("\n", start);
|
|
while (newline !== -1) {
|
|
const line = buffer.slice(start, newline).replace(/\r$/, "").trim();
|
|
if (line) yield JSON.parse(line) as T;
|
|
start = newline + 1;
|
|
newline = buffer.indexOf("\n", start);
|
|
}
|
|
buffer = buffer.slice(start);
|
|
if (done) break;
|
|
}
|
|
const finalLine = buffer.replace(/\r$/, "").trim();
|
|
if (finalLine) yield JSON.parse(finalLine) as T;
|
|
} finally {
|
|
try {
|
|
await reader.cancel();
|
|
} catch {
|
|
// Cancellation can race with an upstream abort.
|
|
}
|
|
reader.releaseLock();
|
|
}
|
|
} finally {
|
|
deadline.clear();
|
|
}
|
|
}
|
|
|
|
/** Raised only when a server does not offer the SSE build-events representation. */
|
|
export class ApiStreamUnsupportedError extends Error {}
|
|
|
|
/** One authenticated SSE connection. The caller owns reconnection and its cursor. */
|
|
export async function* apiStreamEvents<T>(
|
|
path: string,
|
|
after: number,
|
|
signal?: AbortSignal,
|
|
options: ApiRequestOptions = {},
|
|
): AsyncGenerator<{ event: T; id?: number; type?: string }, void, void> {
|
|
const deadline = requestDeadline(signal, 0);
|
|
const headers = new Headers({ accept: "text/event-stream" });
|
|
headers.set("last-event-id", String(after));
|
|
try {
|
|
const response = await sendRequest(
|
|
path,
|
|
{ headers },
|
|
options,
|
|
deadline.signal,
|
|
);
|
|
if ([404, 405, 406, 415, 501].includes(response.status))
|
|
throw new ApiStreamUnsupportedError("Build SSE events are unavailable");
|
|
await assertResponseOk(response);
|
|
if (
|
|
!response.headers
|
|
.get("content-type")
|
|
?.toLowerCase()
|
|
.startsWith("text/event-stream")
|
|
) {
|
|
await response.body?.cancel();
|
|
throw new ApiStreamUnsupportedError("Build SSE events are unavailable");
|
|
}
|
|
if (!response.body) throw new Error("Empty build SSE response");
|
|
const reader = response.body.getReader();
|
|
const decoder = new TextDecoder();
|
|
let buffer = "";
|
|
let data: string[] = [];
|
|
let eventType = "";
|
|
let id: number | undefined;
|
|
try {
|
|
for (;;) {
|
|
const { value, done } = await reader.read();
|
|
buffer += decoder.decode(value, { stream: !done });
|
|
while (buffer.includes("\n")) {
|
|
const newline = buffer.indexOf("\n");
|
|
const line = buffer.slice(0, newline).replace(/\r$/, "");
|
|
buffer = buffer.slice(newline + 1);
|
|
if (line === "") {
|
|
if (
|
|
data.length &&
|
|
(eventType === "log" ||
|
|
eventType === "status" ||
|
|
eventType === "gap")
|
|
)
|
|
yield {
|
|
event: JSON.parse(data.join("\n")) as T,
|
|
id,
|
|
...(eventType === "gap" && { type: eventType }),
|
|
};
|
|
data = [];
|
|
id = undefined;
|
|
eventType = "";
|
|
} else if (line.startsWith("data:"))
|
|
data.push(line.slice(5).replace(/^ /, ""));
|
|
else if (line.startsWith("event:")) eventType = line.slice(6).trim();
|
|
else if (line.startsWith("id:")) {
|
|
const raw = line.slice(3).trim();
|
|
if (/^(0|[1-9]\d*)$/.test(raw) && Number.isSafeInteger(Number(raw)))
|
|
id = Number(raw);
|
|
}
|
|
}
|
|
// Bound a malformed or non-SSE response that never terminates a line.
|
|
if (buffer.length > 2 * 1024 * 1024)
|
|
throw new Error("Build SSE line too long");
|
|
if (done) break;
|
|
}
|
|
} finally {
|
|
try {
|
|
await reader.cancel();
|
|
} catch {
|
|
/* Upstream aborted. */
|
|
}
|
|
reader.releaseLock();
|
|
}
|
|
} finally {
|
|
deadline.clear();
|
|
}
|
|
}
|