287 lines
7.7 KiB
TypeScript
287 lines
7.7 KiB
TypeScript
import { AsyncLocalStorage } from "node:async_hooks";
|
|
|
|
export type ProcessLogEntry = Record<string, unknown>;
|
|
|
|
export type ProcessLogger = {
|
|
log(entry: ProcessLogEntry): void;
|
|
};
|
|
|
|
const MAX_RESPONSE_LOG_BYTES = 64 * 1024;
|
|
const requestContext = new AsyncLocalStorage<{ active: boolean }>();
|
|
|
|
export function logValue(
|
|
value: unknown,
|
|
seen = new WeakSet<object>(),
|
|
): unknown {
|
|
if (
|
|
value === null ||
|
|
typeof value === "string" ||
|
|
typeof value === "number" ||
|
|
typeof value === "boolean"
|
|
)
|
|
return value;
|
|
if (typeof value === "undefined") return undefined;
|
|
if (typeof value === "bigint") return value.toString();
|
|
if (typeof value === "symbol" || typeof value === "function")
|
|
return String(value);
|
|
if (value instanceof Uint8Array)
|
|
return { encoding: "base64", data: Buffer.from(value).toString("base64") };
|
|
if (value instanceof Error)
|
|
return {
|
|
name: value.name,
|
|
message: value.message,
|
|
...(value.stack && { stack: value.stack }),
|
|
};
|
|
if (typeof value !== "object") return String(value);
|
|
if (seen.has(value)) return "[Circular]";
|
|
seen.add(value);
|
|
if (Array.isArray(value)) return value.map((item) => logValue(item, seen));
|
|
try {
|
|
return Object.fromEntries(
|
|
Object.entries(value).map(([key, item]) => [key, logValue(item, seen)]),
|
|
);
|
|
} catch (error) {
|
|
return { unserializable: String(error) };
|
|
}
|
|
}
|
|
|
|
const METHOD_LABELS: Record<string, string> = {
|
|
GET: "GET",
|
|
POST: "PST",
|
|
PUT: "PUT",
|
|
PATCH: "PTC",
|
|
DELETE: "DEL",
|
|
HEAD: "HED",
|
|
OPTIONS: "OPT",
|
|
CONNECT: "CON",
|
|
TRACE: "TRC",
|
|
};
|
|
|
|
function requestLine(entry: ProcessLogEntry): string | undefined {
|
|
const event = entry.event;
|
|
if (
|
|
event === "kuber.server.request.start" ||
|
|
event === "kuber.k8s.request.start" ||
|
|
event === "kuber.server.response.body"
|
|
)
|
|
return "";
|
|
const kubernetes =
|
|
event === "kuber.k8s.request.end" ||
|
|
event === "kuber.k8s.request.failed";
|
|
if (kubernetes && !requestContext.getStore()?.active) return "";
|
|
if (
|
|
!kubernetes &&
|
|
event !== "kuber.server.request.end" &&
|
|
event !== "kuber.server.request.failed"
|
|
)
|
|
return;
|
|
if (typeof entry.method !== "string" || typeof entry.pathname !== "string")
|
|
return;
|
|
if (
|
|
!kubernetes &&
|
|
entry.method.toUpperCase() === "GET" &&
|
|
entry.pathname === "/api/v2/health"
|
|
)
|
|
return "";
|
|
|
|
const method = entry.method.toUpperCase();
|
|
const label = METHOD_LABELS[method] ?? method.slice(0, 3).padEnd(3, "_");
|
|
const status =
|
|
event === "kuber.k8s.request.failed" ||
|
|
event === "kuber.server.request.failed"
|
|
? "ERR"
|
|
: entry.status;
|
|
if (typeof status !== "number" && status !== "ERR") return;
|
|
const indent = kubernetes ? " " : "";
|
|
return `${indent}${label} ${entry.pathname} ${status}`;
|
|
}
|
|
|
|
export const processLogger: ProcessLogger = {
|
|
log(entry) {
|
|
try {
|
|
const line = requestLine(entry);
|
|
if (line) console.log(line);
|
|
} catch {
|
|
// Process logging must never change request behavior.
|
|
}
|
|
},
|
|
};
|
|
|
|
export function safeLog(logger: ProcessLogger, entry: ProcessLogEntry): void {
|
|
try {
|
|
logger.log(entry);
|
|
} catch {
|
|
// Logging is deliberately best-effort and process-only.
|
|
}
|
|
}
|
|
|
|
function headers(headers: Headers): Record<string, string> {
|
|
return Object.fromEntries(headers);
|
|
}
|
|
|
|
function isStreamingContentType(value: string | null): boolean {
|
|
const contentType = value?.split(";", 1)[0]?.trim().toLowerCase();
|
|
return (
|
|
contentType === "text/event-stream" ||
|
|
contentType === "application/x-ndjson" ||
|
|
contentType === "application/ndjson" ||
|
|
contentType === "application/json-seq"
|
|
);
|
|
}
|
|
|
|
function responseBodySize(value: Response): number | undefined {
|
|
const contentLength = value.headers.get("content-length");
|
|
if (!contentLength || !/^\d+$/.test(contentLength)) return;
|
|
const size = Number(contentLength);
|
|
return Number.isSafeInteger(size) ? size : undefined;
|
|
}
|
|
|
|
async function readBoundedResponseBody(
|
|
body: ReadableStream<Uint8Array>,
|
|
): Promise<Uint8Array> {
|
|
const reader = body.getReader();
|
|
const chunks: Uint8Array[] = [];
|
|
let size = 0;
|
|
for (;;) {
|
|
const { done, value } = await reader.read();
|
|
if (done) break;
|
|
size += value.byteLength;
|
|
if (size > MAX_RESPONSE_LOG_BYTES) {
|
|
void reader.cancel();
|
|
throw new RangeError("Response log body exceeds limit");
|
|
}
|
|
chunks.push(value);
|
|
}
|
|
const result = new Uint8Array(size);
|
|
let offset = 0;
|
|
for (const chunk of chunks) {
|
|
result.set(chunk, offset);
|
|
offset += chunk.byteLength;
|
|
}
|
|
return result;
|
|
}
|
|
|
|
function logResponseBody(
|
|
value: Response,
|
|
logger: ProcessLogger,
|
|
entry: ProcessLogEntry,
|
|
): void {
|
|
if (!value.body) return;
|
|
|
|
if (isStreamingContentType(value.headers.get("content-type"))) {
|
|
safeLog(logger, { ...entry, bodySkipped: true, streaming: true });
|
|
return;
|
|
}
|
|
|
|
const size = responseBodySize(value);
|
|
if (size === undefined) {
|
|
safeLog(logger, { ...entry, bodySkipped: true, bodySizeUnknown: true });
|
|
return;
|
|
}
|
|
if (size > MAX_RESPONSE_LOG_BYTES) {
|
|
safeLog(logger, {
|
|
...entry,
|
|
bodySkipped: true,
|
|
bodyTooLarge: true,
|
|
contentLength: size,
|
|
});
|
|
return;
|
|
}
|
|
|
|
try {
|
|
const body = value.clone().body;
|
|
if (!body) return;
|
|
void readBoundedResponseBody(body)
|
|
.then((body) => safeLog(logger, { ...entry, body: logValue(body) }))
|
|
.catch((error) =>
|
|
safeLog(
|
|
logger,
|
|
error instanceof RangeError
|
|
? { ...entry, bodySkipped: true, bodyTooLarge: true }
|
|
: { ...entry, error: logValue(error) },
|
|
),
|
|
);
|
|
} catch (error) {
|
|
safeLog(logger, { ...entry, error: logValue(error) });
|
|
}
|
|
}
|
|
|
|
export async function logServerRequest<T>(
|
|
request: Request,
|
|
requestId: string,
|
|
handler: () => Promise<T> | T,
|
|
logger: ProcessLogger = processLogger,
|
|
): Promise<T> {
|
|
const startedAt = Date.now();
|
|
const url = new URL(request.url);
|
|
const base = {
|
|
type: "server.request",
|
|
requestId,
|
|
method: request.method,
|
|
url: request.url,
|
|
pathname: url.pathname,
|
|
headers: headers(request.headers),
|
|
};
|
|
safeLog(logger, { event: "kuber.server.request.start", ...base });
|
|
const context = { active: true };
|
|
return requestContext.run(context, async () => {
|
|
try {
|
|
const response = await handler();
|
|
safeLog(logger, {
|
|
event: "kuber.server.request.end",
|
|
...base,
|
|
status: response instanceof Response ? response.status : 101,
|
|
durationMs: Date.now() - startedAt,
|
|
...(response instanceof Response && {
|
|
responseHeaders: headers(response.headers),
|
|
}),
|
|
});
|
|
if (response instanceof Response)
|
|
logResponseBody(response, logger, {
|
|
event: "kuber.server.response.body",
|
|
...base,
|
|
status: response.status,
|
|
});
|
|
return response;
|
|
} catch (error) {
|
|
safeLog(logger, {
|
|
event: "kuber.server.request.failed",
|
|
...base,
|
|
durationMs: Date.now() - startedAt,
|
|
error: logValue(error),
|
|
});
|
|
throw error;
|
|
} finally {
|
|
context.active = false;
|
|
}
|
|
});
|
|
}
|
|
|
|
export async function logKubernetesRequest<T>(
|
|
request: Omit<ProcessLogEntry, "event" | "type">,
|
|
handler: () => Promise<T>,
|
|
logger: ProcessLogger = processLogger,
|
|
): Promise<T> {
|
|
const startedAt = Date.now();
|
|
const base = { type: "kubernetes.request", ...request };
|
|
safeLog(logger, { event: "kuber.k8s.request.start", ...base });
|
|
try {
|
|
const result = await handler();
|
|
safeLog(logger, {
|
|
event: "kuber.k8s.request.end",
|
|
...base,
|
|
status: 200,
|
|
durationMs: Date.now() - startedAt,
|
|
});
|
|
return result;
|
|
} catch (error) {
|
|
safeLog(logger, {
|
|
event: "kuber.k8s.request.failed",
|
|
...base,
|
|
durationMs: Date.now() - startedAt,
|
|
error: logValue(error),
|
|
});
|
|
throw error;
|
|
}
|
|
}
|