84 lines
3.2 KiB
TypeScript
84 lines
3.2 KiB
TypeScript
import type { LogEntryBody, LogStreamEventResponse, LogStreamEventsSessionResponse, LogStreamResponse } from "../api/types";
|
|
|
|
export type LiveLogEntry = LogEntryBody & { source: string; streamId: string; streamKey: string };
|
|
|
|
export function parseLogStreamEvent(event: MessageEvent): LogStreamResponse | null {
|
|
const value = parseEventData(event);
|
|
if (!isRecord(value) || typeof value.id !== "string" || typeof value.serverInstanceId !== "string") return null;
|
|
return value as unknown as LogStreamResponse;
|
|
}
|
|
|
|
export function parseServerLogEvent(event: MessageEvent): LogStreamEventResponse | null {
|
|
const value = parseEventData(event);
|
|
if (!isRecord(value) || typeof value.streamId !== "string" || !isRecord(value.entry)) return null;
|
|
return value as unknown as LogStreamEventResponse;
|
|
}
|
|
|
|
export function parseLogSessionEvent(event: MessageEvent): LogStreamEventsSessionResponse | null {
|
|
const value = parseEventData(event);
|
|
if (!isRecord(value) || typeof value.serverInstanceId !== "string") return null;
|
|
return value as unknown as LogStreamEventsSessionResponse;
|
|
}
|
|
|
|
export function entryFromServerLogEvent(event: LogStreamEventResponse): LiveLogEntry {
|
|
return { ...event.entry, source: event.source || event.streamKey, streamId: event.streamId, streamKey: event.streamKey };
|
|
}
|
|
|
|
export function streamFromServerLogEvent(event: LogStreamEventResponse): LogStreamResponse {
|
|
return {
|
|
id: event.streamId,
|
|
serverInstanceId: event.serverInstanceId,
|
|
source: event.source,
|
|
streamKey: event.streamKey,
|
|
logSessionId: event.logSessionId,
|
|
sessionStartedAt: event.sessionStartedAt,
|
|
latestSeq: event.latestSeq,
|
|
storageBackend: "",
|
|
retentionPolicy: "",
|
|
createdAt: "",
|
|
updatedAt: new Date().toISOString()
|
|
};
|
|
}
|
|
|
|
export function mergeLogStreams(current: LogStreamResponse[], incoming: LogStreamResponse): LogStreamResponse[] {
|
|
const index = current.findIndex((stream) => stream.id === incoming.id);
|
|
if (index === -1) return [...current, incoming].sort(compareLogStreams);
|
|
const next = [...current];
|
|
next[index] = { ...next[index], ...incoming, latestSeq: Math.max(next[index].latestSeq, incoming.latestSeq) };
|
|
return next.sort(compareLogStreams);
|
|
}
|
|
|
|
export function appendLiveLogEntries(current: LiveLogEntry[], incoming: LiveLogEntry[], limit: number): LiveLogEntry[] {
|
|
const seen = new Set(current.map(logEntryKey));
|
|
const next = [...current];
|
|
for (const entry of incoming) {
|
|
const key = logEntryKey(entry);
|
|
if (seen.has(key)) continue;
|
|
seen.add(key);
|
|
next.push(entry);
|
|
}
|
|
return next.slice(-limit);
|
|
}
|
|
|
|
function compareLogStreams(a: LogStreamResponse, b: LogStreamResponse): number {
|
|
const updated = (Date.parse(b.updatedAt) || 0) - (Date.parse(a.updatedAt) || 0);
|
|
if (updated !== 0) return updated;
|
|
return b.latestSeq - a.latestSeq || a.streamKey.localeCompare(b.streamKey) || a.id.localeCompare(b.id);
|
|
}
|
|
|
|
function logEntryKey(entry: LiveLogEntry): string {
|
|
return `${entry.streamId}:${entry.seq}`;
|
|
}
|
|
|
|
function parseEventData(event: MessageEvent): unknown {
|
|
try {
|
|
return JSON.parse(String(event.data));
|
|
} catch {
|
|
return null;
|
|
}
|
|
}
|
|
|
|
function isRecord(value: unknown): value is Record<string, unknown> {
|
|
return Boolean(value && typeof value === "object" && !Array.isArray(value));
|
|
}
|