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>; }; export type KubernetesService = { name: string; selector: Readonly>; }; 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>, signal?: AbortSignal, ): Promise; getService( namespace: string, name: string, signal?: AbortSignal, ): Promise; listPods( namespace: string, selector: Readonly>, signal?: AbortSignal, ): Promise; readContainerLogs( request: ContainerLogRequest, signal?: AbortSignal, ): Promise; streamContainerLogs( request: ContainerLogRequest, signal: AbortSignal, ): AsyncIterable; } 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 & { sinceTime?: string; timestamps: boolean; }; type QueueWaiter = (value: T | undefined) => void; class BoundedAsyncQueue { private readonly items: T[] = []; private readonly readers: QueueWaiter[] = []; private readonly writers: Array<() => void> = []; private closed = false; constructor(private readonly capacity: number) {} async push(value: T): Promise { 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((resolve) => this.writers.push(resolve)); } return false; } async shift(): Promise { const item = this.items.shift(); if (item !== undefined) { this.writers.shift()?.(); return item; } if (this.closed) return; return new Promise((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 { if ( error && typeof error === "object" && "code" in error && typeof error.code === "string" ) return false; 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, ): Promise { 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, source: LogContainer, timestamps: boolean, now: () => Date, signal: AbortSignal, emit: (event: LogLineEvent) => Promise, ): Promise { 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 { 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 { const options = normalizeOptions(input); const targets: Array<{ kind: "deployment" | "service"; name: string; selector: Readonly>; }> = []; 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(); 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 { 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 { const options = normalizeOptions(input); const capacity = positiveInteger(input.queueCapacity, 128); const discoveryInterval = positiveInteger(input.discoveryIntervalMs, 2_000); const heartbeatInterval = positiveInteger( input.heartbeatIntervalMs, 5_000, ); const retryInterval = positiveInteger(input.retryIntervalMs, 1_000); const queue = new BoundedAsyncQueue(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(); const completed = new Set(); const retryNotBefore = new Map(); const streamTasks = new Set>(); const pushError = async ( error: unknown, source?: LogContainer, ): Promise => { 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; 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) throw new KubernetesLogError( "Kubernetes log stream ended unexpectedly", true, ); } 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 => { 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 => { 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); }