import * as vscode from 'vscode'; import { LogEntry, StreamInfo, LogFilter } from './types'; export class MultiStreamService { private streams: Map = new Map(); private logsMap: Map = new Map(); private filtersMap: Map = new Map(); private pausedLogsCache: Map = new Map(); private maxLogsPerStream = 20000; private activeStreamId: string | undefined = undefined; private pendingLogUpdates: Map = new Map(); private refreshTimer: NodeJS.Timeout | undefined = undefined; private readonly refreshDebounceMs = 150; private readonly _onDidChangeLog = new vscode.EventEmitter(); public readonly onDidChangeLog = this._onDidChangeLog.event; private readonly _onDidChangeStreams = new vscode.EventEmitter(); public readonly onDidChangeStreams = this._onDidChangeStreams.event; private readonly _onDidChangeActiveStream = new vscode.EventEmitter(); public readonly onDidChangeActiveStream = this._onDidChangeActiveStream.event; private scheduleRefresh(streamId: string): void { this.pendingLogUpdates.set(streamId, true); if (this.refreshTimer) { return; } this.refreshTimer = setTimeout(() => { this.refreshTimer = undefined; if (this.pendingLogUpdates.size > 0) { this.pendingLogUpdates.clear(); this._onDidChangeLog.fire(); } }, this.refreshDebounceMs); } addLogEntry(streamId: string, streamName: string, log: LogEntry): void { if (!this.streams.has(streamId)) { const streamInfo: StreamInfo = { id: streamId, name: streamName, enabled: true }; this.streams.set(streamId, streamInfo); this.logsMap.set(streamId, []); this.filtersMap.set(streamId, {}); this.pausedLogsCache.set(streamId, []); this._onDidChangeStreams.fire(); this.activeStreamId = streamId; this._onDidChangeActiveStream.fire(); } if (!this.streams.get(streamId)?.enabled) { const cachedLogs = this.pausedLogsCache.get(streamId) || []; cachedLogs.push(log); this.pausedLogsCache.set(streamId, cachedLogs); return; } this.handleLogEntry(streamId, streamName, log); } async removeStream(streamId: string): Promise { if (this.streams.has(streamId)) { this.streams.delete(streamId); this.logsMap.delete(streamId); this.filtersMap.delete(streamId); this.pausedLogsCache.delete(streamId); if (this.activeStreamId === streamId) { const remainingStreams = Array.from(this.streams.keys()); this.activeStreamId = remainingStreams.length > 0 ? remainingStreams[0] : undefined; this._onDidChangeActiveStream.fire(); } this._onDidChangeStreams.fire(); this._onDidChangeLog.fire(); } } async toggleStream(streamId: string, enabled: boolean): Promise { const streamInfo = this.streams.get(streamId); if (streamInfo) { const wasDisabled = !streamInfo.enabled; streamInfo.enabled = enabled; if (wasDisabled && enabled) { const cachedLogs = this.pausedLogsCache.get(streamId) || []; if (cachedLogs.length > 0) { this.handleLogEntriesBatch(streamId, streamInfo.name, cachedLogs); this.pausedLogsCache.set(streamId, []); } } this._onDidChangeStreams.fire(); } } private handleLogEntry(streamId: string, streamName: string, log: LogEntry): void { log.streamId = streamId; log.streamName = streamName; const logs = this.logsMap.get(streamId) || []; logs.push(log); if (logs.length > this.maxLogsPerStream) { logs.shift(); } this.logsMap.set(streamId, logs); this.scheduleRefresh(streamId); } private handleLogEntriesBatch(streamId: string, streamName: string, logEntries: LogEntry[]): void { if (logEntries.length === 0) { return; } const logs = this.logsMap.get(streamId) || []; for (const log of logEntries) { log.streamId = streamId; log.streamName = streamName; logs.push(log); } while (logs.length > this.maxLogsPerStream) { logs.shift(); } this.logsMap.set(streamId, logs); this._onDidChangeLog.fire(); } getStreams(): StreamInfo[] { return Array.from(this.streams.values()); } getLogsByStream(streamId: string): LogEntry[] { return this.logsMap.get(streamId) || []; } clearAllLogs(streamId: string): void { if (this.logsMap.has(streamId)) { this.logsMap.set(streamId, []); this._onDidChangeLog.fire(); } } getActiveStreamId(): string | undefined { return this.activeStreamId; } setActiveStreamId(streamId: string | undefined): void { if (streamId === undefined || this.streams.has(streamId)) { this.activeStreamId = streamId; this._onDidChangeActiveStream.fire(); } } getFilter(streamId: string): LogFilter { return this.filtersMap.get(streamId) || {}; } setFilter(streamId: string, filter: LogFilter): void { this.filtersMap.set(streamId, filter); this._onDidChangeLog.fire(); } clearFilter(streamId: string): void { this.filtersMap.set(streamId, {}); this._onDidChangeLog.fire(); } dispose(): void { if (this.refreshTimer) { clearTimeout(this.refreshTimer); this.refreshTimer = undefined; } } }