Files
kuber/server/log-service.ts
T
2026-09-03 11:28:30 +07:00

576 lines
16 KiB
TypeScript

export const MANAGED_BY_SELECTOR = {
"app.kubernetes.io/managed-by": "kuber",
} as const;
export type LogTarget =
| { kind: "managed-deployments" }
| { kind: "service"; name: string };
export type KubernetesDeployment = {
name: string;
selector: Readonly<Record<string, string>>;
};
export type KubernetesService = {
name: string;
selector: Readonly<Record<string, string>>;
};
export type KubernetesPod = {
name: string;
uid: string;
phase?: string;
containers: readonly string[];
};
export type ContainerLogRequest = {
namespace: string;
pod: string;
container: string;
tailLines?: number;
sinceSeconds?: number;
sinceTime?: string;
timestamps: boolean;
};
/** The only cluster operations the log service can perform. */
export interface KubernetesLogsBackend {
listDeployments(
namespace: string,
labels: Readonly<Record<string, string>>,
signal?: AbortSignal,
): Promise<readonly KubernetesDeployment[]>;
getService(
namespace: string,
name: string,
signal?: AbortSignal,
): Promise<KubernetesService | undefined>;
listPods(
namespace: string,
selector: Readonly<Record<string, string>>,
signal?: AbortSignal,
): Promise<readonly KubernetesPod[]>;
readContainerLogs(
request: ContainerLogRequest,
signal?: AbortSignal,
): Promise<string>;
streamContainerLogs(
request: ContainerLogRequest,
signal: AbortSignal,
): AsyncIterable<string | Uint8Array>;
}
export type LogOptions = {
namespace: string;
target: LogTarget;
tailLines?: number;
sinceSeconds?: number;
sinceTime?: string | Date;
timestamps?: boolean;
};
export type FollowLogOptions = LogOptions & {
signal: AbortSignal;
queueCapacity?: number;
discoveryIntervalMs?: number;
heartbeatIntervalMs?: number;
retryIntervalMs?: number;
};
export type LogContainer = {
namespace: string;
targetKind: "deployment" | "service";
targetName: string;
pod: string;
podUid: string;
container: string;
};
export type LogLineEvent = LogContainer & {
type: "log";
timestamp: string;
logTimestamp?: string;
message: string;
};
export type LogHeartbeatEvent = {
type: "heartbeat";
timestamp: string;
};
export type LogErrorEvent = {
type: "error";
timestamp: string;
message: string;
retryable: boolean;
retryAfterMs?: number;
namespace: string;
pod?: string;
container?: string;
};
export type LogEvent = LogLineEvent | LogHeartbeatEvent | LogErrorEvent;
export class KubernetesLogError extends Error {
constructor(
message: string,
readonly retryable: boolean,
) {
super(message);
this.name = "KubernetesLogError";
}
}
type NormalizedOptions = Omit<LogOptions, "sinceTime" | "timestamps"> & {
sinceTime?: string;
timestamps: boolean;
};
type QueueWaiter<T> = (value: T | undefined) => void;
class BoundedAsyncQueue<T> {
private readonly items: T[] = [];
private readonly readers: QueueWaiter<T>[] = [];
private readonly writers: Array<() => void> = [];
private closed = false;
constructor(private readonly capacity: number) {}
async push(value: T): Promise<boolean> {
while (!this.closed) {
const reader = this.readers.shift();
if (reader) {
reader(value);
return true;
}
if (this.items.length < this.capacity) {
this.items.push(value);
return true;
}
await new Promise<void>((resolve) => this.writers.push(resolve));
}
return false;
}
async shift(): Promise<T | undefined> {
const item = this.items.shift();
if (item !== undefined) {
this.writers.shift()?.();
return item;
}
if (this.closed) return;
return new Promise<T | undefined>((resolve) => this.readers.push(resolve));
}
close(discard = false): void {
if (this.closed) return;
this.closed = true;
if (discard) this.items.length = 0;
for (const reader of this.readers.splice(0)) reader(undefined);
for (const writer of this.writers.splice(0)) writer();
}
}
function positiveInteger(value: number | undefined, fallback: number): number {
const result = value ?? fallback;
if (!Number.isSafeInteger(result) || result <= 0)
throw new Error(
"Log service intervals and queue capacity must be positive integers",
);
return result;
}
function normalizeOptions(options: LogOptions): NormalizedOptions {
if (!options.namespace.trim()) throw new Error("Namespace is required");
if (options.target.kind === "service" && !options.target.name.trim())
throw new Error("Service name is required");
if (
options.tailLines !== undefined &&
(!Number.isSafeInteger(options.tailLines) || options.tailLines < 0)
)
throw new Error("tailLines must be a non-negative integer");
if (
options.sinceSeconds !== undefined &&
(!Number.isSafeInteger(options.sinceSeconds) || options.sinceSeconds < 0)
)
throw new Error("sinceSeconds must be a non-negative integer");
if (options.sinceSeconds !== undefined && options.sinceTime !== undefined)
throw new Error("sinceSeconds and sinceTime are mutually exclusive");
let sinceTime: string | undefined;
if (options.sinceTime !== undefined) {
const date =
options.sinceTime instanceof Date
? options.sinceTime
: new Date(options.sinceTime);
if (Number.isNaN(date.getTime()))
throw new Error("sinceTime must be valid");
sinceTime = date.toISOString();
}
return { ...options, sinceTime, timestamps: options.timestamps ?? false };
}
function errorMessage(error: unknown): string {
return error instanceof Error
? error.message
: "Kubernetes log request failed";
}
function isAbort(signal: AbortSignal): boolean {
return signal.aborted;
}
function isRetryable(error: unknown): boolean {
return !(error instanceof KubernetesLogError) || error.retryable;
}
function requestFor(
source: LogContainer,
options: NormalizedOptions,
): ContainerLogRequest {
return {
namespace: source.namespace,
pod: source.pod,
container: source.container,
tailLines: options.tailLines,
sinceSeconds: options.sinceSeconds,
sinceTime: options.sinceTime,
timestamps: options.timestamps,
};
}
function parseLine(
line: string,
source: LogContainer,
timestamps: boolean,
now: () => Date,
): LogLineEvent {
let message = line.replace(/\r$/, "");
let logTimestamp: string | undefined;
if (timestamps) {
const separator = message.indexOf(" ");
if (separator > 0) {
const candidate = message.slice(0, separator);
if (!Number.isNaN(Date.parse(candidate))) {
logTimestamp = candidate;
message = message.slice(separator + 1);
}
}
}
return {
type: "log",
timestamp: now().toISOString(),
...(logTimestamp ? { logTimestamp } : {}),
...source,
message,
};
}
async function emitText(
text: string,
source: LogContainer,
timestamps: boolean,
now: () => Date,
emit: (event: LogLineEvent) => void | Promise<unknown>,
): Promise<void> {
if (!text) return;
const lines = text.split(/\n/);
if (text.endsWith("\n")) lines.pop();
for (const line of lines)
await emit(parseLine(line, source, timestamps, now));
}
async function emitChunks(
chunks: AsyncIterable<string | Uint8Array>,
source: LogContainer,
timestamps: boolean,
now: () => Date,
signal: AbortSignal,
emit: (event: LogLineEvent) => Promise<unknown>,
): Promise<void> {
const decoder = new TextDecoder();
let buffer = "";
for await (const chunk of chunks) {
buffer +=
typeof chunk === "string"
? chunk
: decoder.decode(chunk, { stream: true });
let newline = buffer.indexOf("\n");
while (newline >= 0) {
await emit(parseLine(buffer.slice(0, newline), source, timestamps, now));
buffer = buffer.slice(newline + 1);
newline = buffer.indexOf("\n");
}
}
buffer += decoder.decode();
if (buffer && !signal.aborted)
await emit(parseLine(buffer, source, timestamps, now));
}
function wait(ms: number, signal: AbortSignal): Promise<void> {
if (signal.aborted) return Promise.resolve();
return new Promise((resolve) => {
const timer = setTimeout(done, ms);
function done() {
clearTimeout(timer);
signal.removeEventListener("abort", done);
resolve();
}
signal.addEventListener("abort", done, { once: true });
});
}
export function encodeLogEvent(event: LogEvent): string {
return JSON.stringify(event);
}
export class LogService {
constructor(
private readonly backend: KubernetesLogsBackend,
private readonly now: () => Date = () => new Date(),
) {}
async listContainers(
input: LogOptions,
signal?: AbortSignal,
): Promise<LogContainer[]> {
const options = normalizeOptions(input);
const targets: Array<{
kind: "deployment" | "service";
name: string;
selector: Readonly<Record<string, string>>;
}> = [];
if (options.target.kind === "managed-deployments") {
const deployments = await this.backend.listDeployments(
options.namespace,
MANAGED_BY_SELECTOR,
signal,
);
for (const deployment of deployments) {
if (Object.keys(deployment.selector).length > 0)
targets.push({ kind: "deployment", ...deployment });
}
} else {
const service = await this.backend.getService(
options.namespace,
options.target.name,
signal,
);
if (!service)
throw new KubernetesLogError(
`Service ${options.target.name} was not found`,
false,
);
if (Object.keys(service.selector).length === 0)
throw new KubernetesLogError(
`Service ${options.target.name} has no pod selector`,
false,
);
targets.push({ kind: "service", ...service });
}
const result = new Map<string, LogContainer>();
for (const target of targets) {
const pods = await this.backend.listPods(
options.namespace,
target.selector,
signal,
);
for (const pod of pods) {
if (!pod.name || !pod.uid) continue;
for (const container of pod.containers) {
if (!container) continue;
const source: LogContainer = {
namespace: options.namespace,
targetKind: target.kind,
targetName: target.name,
pod: pod.name,
podUid: pod.uid,
container,
};
result.set(`${pod.uid}\0${container}`, source);
}
}
}
return [...result.values()];
}
async collect(input: LogOptions, signal?: AbortSignal): Promise<LogEvent[]> {
const options = normalizeOptions(input);
const sources = await this.listContainers(options, signal);
const events: LogEvent[] = [];
for (const source of sources) {
try {
const text = await this.backend.readContainerLogs(
requestFor(source, options),
signal,
);
await emitText(text, source, options.timestamps, this.now, (event) => {
events.push(event);
});
} catch (error) {
if (signal && isAbort(signal)) throw error;
events.push({
type: "error",
timestamp: this.now().toISOString(),
namespace: options.namespace,
pod: source.pod,
container: source.container,
message: errorMessage(error),
retryable: isRetryable(error),
});
}
}
return events;
}
async *follow(input: FollowLogOptions): AsyncGenerator<LogEvent> {
const options = normalizeOptions(input);
const capacity = positiveInteger(input.queueCapacity, 128);
const discoveryInterval = positiveInteger(input.discoveryIntervalMs, 2_000);
const heartbeatInterval = positiveInteger(
input.heartbeatIntervalMs,
15_000,
);
const retryInterval = positiveInteger(input.retryIntervalMs, 1_000);
const queue = new BoundedAsyncQueue<LogEvent>(capacity);
const controller = new AbortController();
const stop = () => {
queue.close(true);
controller.abort();
};
if (input.signal.aborted) controller.abort();
else input.signal.addEventListener("abort", stop, { once: true });
const active = new Map<string, AbortController>();
const completed = new Set<string>();
const retryNotBefore = new Map<string, number>();
const streamTasks = new Set<Promise<void>>();
const pushError = async (
error: unknown,
source?: LogContainer,
): Promise<void> => {
await queue.push({
type: "error",
timestamp: this.now().toISOString(),
namespace: options.namespace,
pod: source?.pod,
container: source?.container,
message: errorMessage(error),
retryable: isRetryable(error),
...(isRetryable(error) ? { retryAfterMs: retryInterval } : {}),
});
};
const startStream = (source: LogContainer): void => {
const key = `${source.podUid}\0${source.container}`;
const local = new AbortController();
active.set(key, local);
const signal = AbortSignal.any([controller.signal, local.signal]);
let task!: Promise<void>;
task = (async () => {
try {
await emitChunks(
this.backend.streamContainerLogs(
requestFor(source, options),
signal,
),
source,
options.timestamps,
this.now,
signal,
(event) => queue.push(event),
);
if (!signal.aborted) completed.add(key);
} catch (error) {
if (!isAbort(signal)) {
await pushError(error, source);
if (isRetryable(error))
retryNotBefore.set(key, Date.now() + retryInterval);
else completed.add(key);
}
} finally {
active.delete(key);
streamTasks.delete(task);
}
})();
streamTasks.add(task);
};
const discover = async (): Promise<void> => {
while (!controller.signal.aborted) {
try {
const sources = await this.listContainers(options, controller.signal);
const desired = new Set(
sources.map((source) => `${source.podUid}\0${source.container}`),
);
for (const key of completed) {
if (!desired.has(key)) completed.delete(key);
}
for (const key of retryNotBefore.keys()) {
if (!desired.has(key)) retryNotBefore.delete(key);
}
for (const [key, stream] of active) {
if (!desired.has(key)) stream.abort();
}
for (const source of sources) {
const key = `${source.podUid}\0${source.container}`;
if (
!active.has(key) &&
!completed.has(key) &&
(retryNotBefore.get(key) ?? 0) <= Date.now()
) {
retryNotBefore.delete(key);
startStream(source);
}
}
} catch (error) {
if (!isAbort(controller.signal)) {
await pushError(error);
if (!isRetryable(error)) controller.abort();
}
}
await wait(discoveryInterval, controller.signal);
}
};
const heartbeat = async (): Promise<void> => {
while (!controller.signal.aborted) {
await wait(heartbeatInterval, controller.signal);
if (!controller.signal.aborted)
await queue.push({
type: "heartbeat",
timestamp: this.now().toISOString(),
});
}
};
const runner = Promise.all([discover(), heartbeat()]).finally(async () => {
for (const stream of active.values()) stream.abort();
await Promise.allSettled([...streamTasks]);
queue.close();
});
try {
for (;;) {
const event = await queue.shift();
if (event === undefined) return;
yield event;
}
} finally {
input.signal.removeEventListener("abort", stop);
controller.abort();
queue.close(true);
await runner;
}
}
}
export function createLogService(
backend: KubernetesLogsBackend,
now?: () => Date,
): LogService {
return new LogService(backend, now);
}