feat: v2
This commit is contained in:
@@ -0,0 +1,575 @@
|
||||
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);
|
||||
}
|
||||
Reference in New Issue
Block a user