From e00a72bb3d8226ce6eefc6480d92d494597a357c Mon Sep 17 00:00:00 2001 From: Brad Groux Date: Wed, 4 Feb 2026 09:56:19 -0600 Subject: [PATCH] perf: Stream telemetry reads, push pagination to service, optimize lookups MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PERF-001: Replace gunzipSync with streaming readline + createGunzip Apply filters during streaming for early rejection PERF-002: Add offset parameter to activity service getActivities() Route uses offset instead of fetching page*limit and slicing PERF-003: BacklogRepository findById() uses file prefix lookup O(n) → O(1) task-metrics velocity uses Set for archived IDs (O(n²) → O(n)) audit-service verifyAuditLog() uses readline streaming Ref: RF-002a audit findings (Medium/Low severity) --- server/src/routes/activity.ts | 12 +- server/src/services/activity-service.ts | 8 +- server/src/services/audit-service.ts | 74 ++++++++--- server/src/services/metrics/task-metrics.ts | 5 +- server/src/services/telemetry-service.ts | 132 ++++++++++++-------- server/src/storage/backlog-repository.ts | 17 ++- 6 files changed, 162 insertions(+), 86 deletions(-) diff --git a/server/src/routes/activity.ts b/server/src/routes/activity.ts index 3edcb099..3c0a90c6 100644 --- a/server/src/routes/activity.ts +++ b/server/src/routes/activity.ts @@ -38,14 +38,12 @@ router.get( // If pagination is requested, use the sendPaginated helper if (page > 0) { const total = await activityService.countActivities(hasFilters ? filters : undefined); - // Fetch enough items to cover the requested page - const fetchLimit = page * limit; - const allMatching = await activityService.getActivities( - fetchLimit, - hasFilters ? filters : undefined + const offset = (page - 1) * limit; + const paged = await activityService.getActivities( + limit, + hasFilters ? filters : undefined, + offset ); - const start = (page - 1) * limit; - const paged = allMatching.slice(start, start + limit); sendPaginated(res, paged, { page, limit, total }); } else { const activities = await activityService.getActivities( diff --git a/server/src/services/activity-service.ts b/server/src/services/activity-service.ts index 8c3252cb..afd5722f 100644 --- a/server/src/services/activity-service.ts +++ b/server/src/services/activity-service.ts @@ -74,7 +74,11 @@ export class ActivityService { } } - async getActivities(limit: number = 50, filters?: ActivityFilters): Promise { + async getActivities( + limit: number = 50, + filters?: ActivityFilters, + offset: number = 0 + ): Promise { let activities = await this.loadAll(); // Apply filters @@ -99,7 +103,7 @@ export class ActivityService { } } - return activities.slice(0, limit); + return activities.slice(offset, offset + limit); } /** diff --git a/server/src/services/audit-service.ts b/server/src/services/audit-service.ts index dded8566..fe24bc0b 100644 --- a/server/src/services/audit-service.ts +++ b/server/src/services/audit-service.ts @@ -8,8 +8,10 @@ * {dataDir}/audit/audit-{YYYY-MM}.log */ import fs from 'fs/promises'; +import { createReadStream } from 'fs'; import path from 'path'; import crypto from 'crypto'; +import readline from 'readline'; import { createLogger } from '../lib/logger.js'; const log = createLogger('audit'); @@ -176,9 +178,9 @@ async function writeEntry(event: AuditEvent): Promise { * Returns a result indicating whether the chain is intact. */ export async function verifyAuditLog(filePath: string): Promise { - let content: string; + // Check if file exists try { - content = await fs.readFile(filePath, 'utf8'); + await fs.access(filePath); } catch (err: unknown) { if ((err as NodeJS.ErrnoException).code === 'ENOENT') { return { valid: true, entries: 0 }; @@ -186,29 +188,61 @@ export async function verifyAuditLog(filePath: string): Promise { throw err; } - const lines = content.trimEnd().split('\n').filter(Boolean); - if (lines.length === 0) { - return { valid: true, entries: 0 }; - } + return new Promise((resolve, reject) => { + const stream = createReadStream(filePath, { encoding: 'utf8' }); + const rl = readline.createInterface({ + input: stream, + crlfDelay: Infinity, + }); - let prevHash = ''; + let prevHash = ''; + let lineIndex = 0; + let totalLines = 0; + let invalidResult: VerifyResult | null = null; - for (let i = 0; i < lines.length; i++) { - let entry: AuditEntry; - try { - entry = JSON.parse(lines[i]) as AuditEntry; - } catch { - return { valid: false, entries: lines.length, firstBroken: i }; - } + rl.on('line', (line) => { + if (invalidResult) return; // Already found an error - if (entry.integrity !== prevHash) { - return { valid: false, entries: lines.length, firstBroken: i }; - } + const trimmed = line.trim(); + if (!trimmed) { + lineIndex++; + return; + } - prevHash = sha256(lines[i]); - } + totalLines++; - return { valid: true, entries: lines.length }; + let entry: AuditEntry; + try { + entry = JSON.parse(trimmed) as AuditEntry; + } catch { + invalidResult = { valid: false, entries: totalLines, firstBroken: lineIndex }; + rl.close(); + stream.destroy(); + return; + } + + if (entry.integrity !== prevHash) { + invalidResult = { valid: false, entries: totalLines, firstBroken: lineIndex }; + rl.close(); + stream.destroy(); + return; + } + + prevHash = sha256(trimmed); + lineIndex++; + }); + + rl.on('close', () => { + if (invalidResult) { + resolve(invalidResult); + } else { + resolve({ valid: true, entries: totalLines }); + } + }); + + rl.on('error', reject); + stream.on('error', reject); + }); } /** diff --git a/server/src/services/metrics/task-metrics.ts b/server/src/services/metrics/task-metrics.ts index 087c9f4e..d89ae5b8 100644 --- a/server/src/services/metrics/task-metrics.ts +++ b/server/src/services/metrics/task-metrics.ts @@ -126,6 +126,9 @@ export async function computeVelocityMetrics( (t) => !project || t.project === project ); + // Pre-compute archived IDs set for O(1) lookup (avoid O(n²) in loop) + const archivedIds = new Set(archivedTasks.map((a) => a.id)); + // Group tasks by sprint const sprintData = new Map< string, @@ -147,7 +150,7 @@ export async function computeVelocityMetrics( data.total++; // Count completed tasks (done or archived) - const isCompleted = task.status === 'done' || archivedTasks.some((a) => a.id === task.id); + const isCompleted = task.status === 'done' || archivedIds.has(task.id); if (isCompleted) { data.completed++; diff --git a/server/src/services/telemetry-service.ts b/server/src/services/telemetry-service.ts index 2101c50b..0e838a13 100644 --- a/server/src/services/telemetry-service.ts +++ b/server/src/services/telemetry-service.ts @@ -1,8 +1,9 @@ import fs from 'fs/promises'; import { createReadStream, createWriteStream } from '../storage/fs-helpers.js'; import path from 'path'; -import { createGzip, gunzipSync } from 'zlib'; +import { createGzip, createGunzip } from 'zlib'; import { pipeline } from 'stream/promises'; +import readline from 'readline'; import { nanoid } from 'nanoid'; import type { TelemetryEvent, @@ -167,44 +168,36 @@ export class TelemetryService { await this.init(); const { type, since, until, taskId, project, limit } = options; + const types = type ? (Array.isArray(type) ? type : [type]) : null; // Determine which files to read based on date range const files = await this.getEventFiles(since, until); - let events: AnyTelemetryEvent[] = []; + const events: AnyTelemetryEvent[] = []; + // Use streaming with early filtering for (const file of files) { - const fileEvents = await this.readEventFile(file); - events.push(...fileEvents); + await this.streamEventFile(file, (event) => { + // Apply filters during streaming (early rejection) + if (types && !types.includes(event.type)) return; + if (since && event.timestamp < since) return; + if (until && event.timestamp > until) return; + if (taskId && event.taskId !== taskId) return; + if (project && event.project !== project) return; + + events.push(event); + + // Note: Can't early-terminate by limit here because we need to sort first. + // However, filtering during streaming reduces memory usage significantly. + }); } - // Apply filters - events = events.filter((event) => { - // Type filter - if (type) { - const types = Array.isArray(type) ? type : [type]; - if (!types.includes(event.type)) return false; - } - - // Time range filters - if (since && event.timestamp < since) return false; - if (until && event.timestamp > until) return false; - - // Task filter - if (taskId && event.taskId !== taskId) return false; - - // Project filter - if (project && event.project !== project) return false; - - return true; - }); - // Sort by timestamp (newest first) events.sort((a, b) => b.timestamp.localeCompare(a.timestamp)); - // Apply limit + // Apply limit after sort if (limit && events.length > limit) { - events = events.slice(0, limit); + return events.slice(0, limit); } return events; @@ -392,39 +385,70 @@ export class TelemetryService { } /** - * Read events from a single file (supports .ndjson and .ndjson.gz) + * Stream events from a single file (supports .ndjson and .ndjson.gz) + * Uses readline for memory-efficient line-by-line processing. + * Calls the callback for each event; return false to stop early. */ - private async readEventFile(filename: string): Promise { + private async streamEventFile( + filename: string, + callback: (event: AnyTelemetryEvent) => boolean | void + ): Promise { const filepath = path.join(this.telemetryDir, filename); const isGzipped = filename.endsWith('.gz'); try { - let content: string; - if (isGzipped) { - const buffer = await fs.readFile(filepath); - content = gunzipSync(buffer).toString('utf-8'); - } else { - content = await fs.readFile(filepath, 'utf-8'); - } - - const lines = content.trim().split('\n').filter(Boolean); - - return lines - .map((line) => { - try { - return JSON.parse(line) as AnyTelemetryEvent; - } catch { - log.error({ err: line }, '[Telemetry] Failed to parse line'); - return null; - } - }) - .filter((e): e is AnyTelemetryEvent => e !== null); - } catch (error: any) { - if (error.code === 'ENOENT') { - return []; - } - throw error; + await fs.access(filepath); + } catch { + return; // File doesn't exist } + + return new Promise((resolve, reject) => { + let stream = createReadStream(filepath); + + if (isGzipped) { + const gunzip = createGunzip(); + stream = stream.pipe(gunzip) as unknown as ReturnType; + } + + const rl = readline.createInterface({ + input: stream as NodeJS.ReadableStream, + crlfDelay: Infinity, + }); + + let stopped = false; + + rl.on('line', (line) => { + if (stopped || !line.trim()) return; + + try { + const event = JSON.parse(line) as AnyTelemetryEvent; + const shouldContinue = callback(event); + if (shouldContinue === false) { + stopped = true; + rl.close(); + stream.destroy(); + } + } catch { + log.error({ err: line }, '[Telemetry] Failed to parse line'); + } + }); + + rl.on('close', () => resolve()); + rl.on('error', reject); + stream.on('error', reject); + }); + } + + /** + * Read all events from a single file (backwards compat wrapper) + * Prefer streamEventFile for large files with filtering/limits. + */ + private async readEventFile(filename: string): Promise { + const events: AnyTelemetryEvent[] = []; + await this.streamEventFile(filename, (event) => { + events.push(event); + }); + return events; } /** diff --git a/server/src/storage/backlog-repository.ts b/server/src/storage/backlog-repository.ts index ead89308..c87eb32b 100644 --- a/server/src/storage/backlog-repository.ts +++ b/server/src/storage/backlog-repository.ts @@ -130,10 +130,23 @@ export class BacklogRepository { /** * Get a single backlog task by ID + * Uses direct file lookup by ID prefix instead of scanning all files. */ async findById(id: string): Promise { - const tasks = await this.listAll(); - return tasks.find((t) => t.id === id) || null; + await this.ensureDirectory(); + + // Files are named: ${id}-${slug}.md + // Find file that starts with the ID + const files = await fs.readdir(this.backlogDir); + const targetFile = files.find((f) => f.startsWith(`${id}-`) && f.endsWith('.md')); + + if (!targetFile) { + return null; + } + + const filepath = path.join(this.backlogDir, targetFile); + const content = await fs.readFile(filepath, 'utf-8'); + return this.parseTaskFile(content, targetFile); } /**