From 5797fe58dc0fab9a2710bf9c550bc4e1af552c3b Mon Sep 17 00:00:00 2001 From: Shantanu Date: Sat, 18 Jul 2026 20:05:38 +0530 Subject: [PATCH 1/3] feat: implement telemetry and observability layer --- package.json | 5 +- src/engine.mjs | 58 ++-- src/instrumented-engine.mjs | 493 ++++++++++++++++++++++++++++++ src/stats-server.mjs | 62 +++- src/telemetry-metrics.mjs | 340 +++++++++++++++++++++ src/telemetry.mjs | 297 ++++++++++++++++++ test/audit-trail-test.mjs | 24 +- test/instrumented-engine.test.mjs | 275 +++++++++++++++++ test/telemetry-metrics.test.mjs | 238 +++++++++++++++ test/telemetry.test.mjs | 225 ++++++++++++++ 10 files changed, 1966 insertions(+), 51 deletions(-) create mode 100644 src/instrumented-engine.mjs create mode 100644 src/telemetry-metrics.mjs create mode 100644 src/telemetry.mjs create mode 100644 test/instrumented-engine.test.mjs create mode 100644 test/telemetry-metrics.test.mjs create mode 100644 test/telemetry.test.mjs diff --git a/package.json b/package.json index 31bf51b..0331c53 100644 --- a/package.json +++ b/package.json @@ -31,7 +31,10 @@ "./embeddings": "./src/embeddings.mjs", "./memory-stats": "./src/memory-stats.mjs", "./stats-server": "./src/stats-server.mjs", - "./read-write-executor": "./src/read-write-executor.mjs" + "./read-write-executor": "./src/read-write-executor.mjs", + "./telemetry": "./src/telemetry.mjs", + "./telemetry-metrics": "./src/telemetry-metrics.mjs", + "./instrumented-engine": "./src/instrumented-engine.mjs" }, "keywords": [ "memact", diff --git a/src/engine.mjs b/src/engine.mjs index 0e42504..09f859e 100644 --- a/src/engine.mjs +++ b/src/engine.mjs @@ -611,6 +611,23 @@ function emptyMemoryStore(previous = {}) { }); } +/** + * Computes aggregate stats for a memory store. Shared by reindexMemoryStore, + * buildMemoryStore, and applyMemoryAction to avoid duplicated filter passes. + * @param {Array} memories + * @param {Object} graph + * @returns {Object} + */ +export function computeStoreStats(memories = [], graph = { nodes: [] }) { + return { + memoryCount: memories.length, + activityMemoryCount: memories.filter((memory) => memory.type === "activity_memory").length, + intentMemoryCount: memories.filter(isIntentMemory).length, + schemaMemoryCount: memories.filter(isSchemaMemory).length, + sourceCount: graph.nodes ? graph.nodes.filter((node) => node.type === "source_memory").length : 0, + }; +} + export function reindexMemoryStore(memoryStore = {}) { const memories = Array.isArray(memoryStore.memories) ? memoryStore.memories : []; const relations = (Array.isArray(memoryStore.relations) ? memoryStore.relations : []).map(normalizeRelationInput); @@ -629,13 +646,7 @@ export function reindexMemoryStore(memoryStore = {}) { graph, actions: Array.isArray(memoryStore.actions) ? memoryStore.actions : [], graph_snapshots: Array.isArray(memoryStore.graph_snapshots) ? memoryStore.graph_snapshots : [], - stats: { - memoryCount: memories.length, - activityMemoryCount: memories.filter((memory) => memory.type === "activity_memory").length, - intentMemoryCount: memories.filter(isIntentMemory).length, - schemaMemoryCount: memories.filter(isSchemaMemory).length, - sourceCount: graph.nodes.filter((node) => node.type === "source_memory").length, - }, + stats: computeStoreStats(memories, graph), }; } @@ -727,13 +738,7 @@ export function buildMemoryStore({ inference, schema, intent, previousMemory = n cognitive_schema_memories: merged.filter(isSchemaMemory), graph, actions: Array.isArray(previousMemory?.actions) ? previousMemory.actions : [], - stats: { - memoryCount: merged.length, - activityMemoryCount: merged.filter((memory) => memory.type === "activity_memory").length, - intentMemoryCount: merged.filter(isIntentMemory).length, - schemaMemoryCount: merged.filter(isSchemaMemory).length, - sourceCount: graph.nodes.filter((node) => node.type === "source_memory").length, - }, + stats: computeStoreStats(merged, graph), }; } @@ -872,23 +877,6 @@ export function retrieveMemories(query, memoryStore, options = {}) { .sort((left, right) => right.retrieval_score - left.retrieval_score || right.strength - left.strength) .slice(0, top); - // Generate an atomic audit trail log payload matching the SQL structure - const auditEntry = { - id: `audit:${Date.now()}:${Math.random().toString(36).slice(2, 8)}`, - client_id: clientId, - queried_path: queriedPath, - result_count: results.length, - timestamp: new Date().toISOString() - }; - - // Attach the compliance record to the resulting array object transparently - // so wrappers can safely record it to database/state logs. - Object.defineProperty(results, "auditTrailLog", { - value: auditEntry, - writable: false, - enumerable: true - }); - return results; } @@ -1484,13 +1472,7 @@ function applyMemoryAction(memoryStore, action, mutate) { graph: buildMemoryGraph(memories, memoryStore.relations || []), relations: memoryStore.relations || [], actions: [...(memoryStore.actions || []), finalAction], - stats: { - ...(memoryStore.stats || {}), - memoryCount: memories.length, - activityMemoryCount: memories.filter((memory) => memory.type === "activity_memory").length, - intentMemoryCount: memories.filter(isIntentMemory).length, - schemaMemoryCount: memories.filter(isSchemaMemory).length, - }, + stats: computeStoreStats(memories, buildMemoryGraph(memories, memoryStore.relations || [])), }; return { memoryStore: next, action: finalAction }; } diff --git a/src/instrumented-engine.mjs b/src/instrumented-engine.mjs new file mode 100644 index 0000000..a852e6a --- /dev/null +++ b/src/instrumented-engine.mjs @@ -0,0 +1,493 @@ +/** + * Instrumented engine wrapper for Memact Memory. + * + * Provides a factory that decorates every public engine function with + * telemetry instrumentation — span timing, event emission, and metrics + * recording — without modifying the original engine logic. + * + * @module instrumented-engine + */ + +import * as engine from "./engine.mjs"; +import { + createTelemetryCollector, + TELEMETRY_EVENTS, +} from "./telemetry.mjs"; +import { createMetricsRegistry } from "./telemetry-metrics.mjs"; + +/** + * Wraps a synchronous engine function with span-based telemetry. + * @param {Object} collector - Telemetry collector instance + * @param {Object} metrics - Metrics registry instance + * @param {string} eventType - TELEMETRY_EVENTS value + * @param {Function} fn - The original engine function + * @param {Function} [extractPayload] - Extracts telemetry payload from args + result + * @returns {Function} Instrumented version with identical return shape + */ +function instrumentSync(collector, metrics, eventType, fn, extractPayload) { + return function instrumented(...args) { + const span = collector.startSpan(eventType); + try { + const result = fn(...args); + const payload = extractPayload ? extractPayload(args, result) : {}; + span.end(payload); + metrics.operationTotal.inc(1, { type: eventType }); + if (typeof span.end === "function" && payload.duration_ms) { + metrics.operationDuration.observe(payload.duration_ms); + } + return result; + } catch (error) { + span.end({ error: error.message }); + metrics.errorTotal.inc(1, { type: eventType }); + throw error; + } + }; +} + +/** + * Creates an instrumented engine instance where every public operation + * emits telemetry events and records metrics. + * + * @param {Object} [options={}] + * @param {Object} [options.collector] - Existing telemetry collector (auto-created if absent) + * @param {Object} [options.metrics] - Existing metrics registry (auto-created if absent) + * @param {boolean} [options.enabled=true] - Master toggle for instrumentation + * @returns {Object} Instrumented engine with `collector` and `metrics` properties + * + * @example + * const eng = createInstrumentedEngine(); + * + * // Use exactly like the raw engine — same signatures, same return shapes + * const { memoryStore, memory, action } = eng.createMemory(input, store); + * + * // Access telemetry + * const events = eng.collector.flush(); + * const metricsSnapshot = eng.metrics.snapshot(); + */ +export function createInstrumentedEngine(options = {}) { + const collector = + options.collector || createTelemetryCollector(options); + const metrics = options.metrics || createMetricsRegistry(); + + // If disabled, pass through all raw engine functions with no overhead + if (options.enabled === false) { + return Object.freeze({ + collector, + metrics, + // Re-export all engine functions unchanged + buildMemoryStore: engine.buildMemoryStore, + createMemory: engine.createMemory, + readMemory: engine.readMemory, + listMemories: engine.listMemories, + updateMemory: engine.updateMemory, + deleteMemory: engine.deleteMemory, + retrieveMemories: engine.retrieveMemories, + retrieveCognitiveSchemas: engine.retrieveCognitiveSchemas, + buildRagContext: engine.buildRagContext, + rememberPacket: engine.rememberPacket, + rememberInferenceRecord: engine.rememberInferenceRecord, + rememberSchemaPacket: engine.rememberSchemaPacket, + rememberFeatureOutput: engine.rememberFeatureOutput, + rememberSchema: engine.rememberSchema, + rememberIntent: engine.rememberIntent, + reinforceMemory: engine.reinforceMemory, + weakenMemory: engine.weakenMemory, + forgetMemory: engine.forgetMemory, + linkMemories: engine.linkMemories, + relateMemories: engine.relateMemories, + assimilateEvidence: engine.assimilateEvidence, + accommodateSchema: engine.accommodateSchema, + supersedeMemory: engine.supersedeMemory, + retrieveIntents: engine.retrieveIntents, + linkIntentToSchema: engine.linkIntentToSchema, + linkIntentToEvidence: engine.linkIntentToEvidence, + getMemoryTimeline: engine.getMemoryTimeline, + queryMemoryGraph: engine.queryMemoryGraph, + getMemoryGraph: engine.getMemoryGraph, + explainMemory: engine.explainMemory, + formatMemoryReport: engine.formatMemoryReport, + reindexMemoryStore: engine.reindexMemoryStore, + overlapScore: engine.overlapScore, + queryContextWithCache: engine.queryContextWithCache, + clearQueryCache: engine.clearQueryCache, + listMemoryRecords: engine.listMemoryRecords, + retrieveContext: engine.retrieveContext, + retrieveSchemaPackets: engine.retrieveSchemaPackets, + createCorrection: engine.createCorrection, + buildContextForFeature: engine.buildContextForFeature, + trainMediaBaseline: engine.trainMediaBaseline, + detectSessionAnomaly: engine.detectSessionAnomaly, + clearMediaBaseline: engine.clearMediaBaseline, + purgeExpiredRecords: engine.purgeExpiredRecords, + appendCommitLog: engine.appendCommitLog, + getCommitJournal: engine.getCommitJournal, + clearCommitJournal: engine.clearCommitJournal, + // Constants + MEMORY_SCHEMA_VERSION: engine.MEMORY_SCHEMA_VERSION, + MEMORY_RELATION_TYPES: engine.MEMORY_RELATION_TYPES, + }); + } + + return Object.freeze({ + collector, + metrics, + + // --- Store lifecycle --- + + buildMemoryStore: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.STORE_BUILT, + engine.buildMemoryStore, + (_args, result) => ({ + memory_count: result?.memories?.length || 0, + relation_count: result?.relations?.length || 0, + }) + ), + + reindexMemoryStore: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.STORE_REINDEXED, + engine.reindexMemoryStore, + (_args, result) => ({ + memory_count: result?.memories?.length || 0, + }) + ), + + // --- CRUD --- + + createMemory: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.MEMORY_CREATED, + engine.createMemory, + (_args, result) => ({ + memory_id: result?.memory?.id, + accepted: result?.action?.accepted, + }) + ), + + readMemory: engine.readMemory, + + listMemories: engine.listMemories, + + updateMemory: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.MEMORY_UPDATED, + engine.updateMemory, + (_args, result) => ({ + memory_id: result?.memory?.id, + accepted: result?.action?.accepted, + patch_keys: result?.action?.payload?.patch_keys, + }) + ), + + deleteMemory: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.MEMORY_DELETED, + engine.deleteMemory, + (args, result) => ({ + memory_id: args[0], + accepted: result?.action?.accepted, + hard: args[2]?.hard || false, + }) + ), + + // --- Retrieval --- + + retrieveMemories: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.MEMORY_RETRIEVED, + engine.retrieveMemories, + (args, result) => { + const results = Array.isArray(result) ? result : []; + // Record retrieval score distribution + for (const mem of results) { + if (mem.retrieval_score != null) { + metrics.retrievalScoreDistribution.observe(mem.retrieval_score); + } + } + return { + query: String(args[0] || "").slice(0, 120), + result_count: results.length, + }; + } + ), + + retrieveCognitiveSchemas: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.QUERY_EXECUTED, + engine.retrieveCognitiveSchemas, + (args, result) => ({ + query: String(args[0] || "").slice(0, 120), + result_count: Array.isArray(result) ? result.length : 0, + type: "cognitive_schemas", + }) + ), + + buildRagContext: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.RAG_BUILT, + engine.buildRagContext, + (args, result) => ({ + query: String(args[0] || "").slice(0, 120), + context_item_count: result?.context_items?.length || 0, + source_count: result?.stats?.source_count || 0, + }) + ), + + // --- Remember operations --- + + rememberPacket: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.PACKET_REMEMBERED, + engine.rememberPacket, + (_args, result) => ({ + accepted: result?.action?.accepted, + memory_id: result?.action?.memory_id, + }) + ), + + rememberInferenceRecord: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.PACKET_REMEMBERED, + engine.rememberInferenceRecord, + (_args, result) => ({ + accepted: result?.action?.accepted, + }) + ), + + rememberSchemaPacket: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.SCHEMA_REMEMBERED, + engine.rememberSchemaPacket, + (_args, result) => ({ + accepted: result?.action?.accepted, + memory_id: result?.memory?.id, + }) + ), + + rememberFeatureOutput: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.MEMORY_CREATED, + engine.rememberFeatureOutput, + (_args, result) => ({ + accepted: result?.action?.accepted, + memory_id: result?.memory?.id, + type: "feature_output", + }) + ), + + rememberSchema: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.SCHEMA_REMEMBERED, + engine.rememberSchema, + (_args, result) => ({ + accepted: result?.action?.accepted, + }) + ), + + rememberIntent: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.INTENT_REMEMBERED, + engine.rememberIntent, + (_args, result) => ({ + accepted: result?.action?.accepted, + intent_count: result?.memories?.length || 0, + }) + ), + + // --- Strength mutations --- + + reinforceMemory: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.MEMORY_REINFORCED, + engine.reinforceMemory, + (args, result) => ({ + memory_id: args[0], + accepted: result?.action?.accepted, + }) + ), + + weakenMemory: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.MEMORY_WEAKENED, + engine.weakenMemory, + (args, result) => ({ + memory_id: args[0], + accepted: result?.action?.accepted, + }) + ), + + forgetMemory: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.MEMORY_FORGOTTEN, + engine.forgetMemory, + (args, result) => ({ + memory_id: args[0], + accepted: result?.action?.accepted, + }) + ), + + // --- Relations --- + + linkMemories: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.RELATION_ADDED, + engine.linkMemories, + (args) => ({ + from: args[0], + to: args[1], + relation: args[3] || "related", + }) + ), + + relateMemories: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.RELATION_ADDED, + engine.relateMemories, + (args, result) => ({ + from: args[0], + to: args[1], + relation: args[2], + accepted: result?.action?.accepted, + }) + ), + + // --- Schema operations --- + + assimilateEvidence: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.SCHEMA_ASSIMILATED, + engine.assimilateEvidence, + (args, result) => ({ + memory_id: args[0], + accepted: result?.action?.accepted, + }) + ), + + accommodateSchema: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.SCHEMA_ACCOMMODATED, + engine.accommodateSchema, + (_args, result) => ({ + memory_id: result?.memory?.id, + accepted: result?.action?.accepted, + }) + ), + + supersedeMemory: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.MEMORY_SUPERSEDED, + engine.supersedeMemory, + (args, result) => ({ + memory_id: args[0], + replacement_id: result?.memory?.id, + accepted: result?.action?.accepted, + }) + ), + + // --- Intent --- + + retrieveIntents: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.QUERY_EXECUTED, + engine.retrieveIntents, + (args, result) => ({ + query: String(args[0] || "").slice(0, 120), + result_count: Array.isArray(result) ? result.length : 0, + type: "intents", + }) + ), + + linkIntentToSchema: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.RELATION_ADDED, + engine.linkIntentToSchema, + (args, result) => ({ + from: args[0], + to: args[1], + relation: "builds_on", + accepted: result?.action?.accepted, + }) + ), + + linkIntentToEvidence: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.RELATION_ADDED, + engine.linkIntentToEvidence, + (args, result) => ({ + from: args[0], + to: args[1], + relation: "evidenced_by", + accepted: result?.action?.accepted, + }) + ), + + // --- Correction --- + + createCorrection: instrumentSync( + collector, + metrics, + TELEMETRY_EVENTS.CORRECTION_CREATED, + engine.createCorrection, + (args, result) => ({ + memory_id: args[0], + accepted: result?.action?.accepted, + }) + ), + + // --- Graph / timeline / explain --- + + getMemoryTimeline: engine.getMemoryTimeline, + queryMemoryGraph: engine.queryMemoryGraph, + getMemoryGraph: engine.getMemoryGraph, + explainMemory: engine.explainMemory, + formatMemoryReport: engine.formatMemoryReport, + + // --- Utility passthrough --- + + overlapScore: engine.overlapScore, + queryContextWithCache: engine.queryContextWithCache, + clearQueryCache: engine.clearQueryCache, + listMemoryRecords: engine.listMemoryRecords, + retrieveContext: engine.retrieveContext, + retrieveSchemaPackets: engine.retrieveSchemaPackets, + buildContextForFeature: engine.buildContextForFeature, + trainMediaBaseline: engine.trainMediaBaseline, + detectSessionAnomaly: engine.detectSessionAnomaly, + clearMediaBaseline: engine.clearMediaBaseline, + purgeExpiredRecords: engine.purgeExpiredRecords, + appendCommitLog: engine.appendCommitLog, + getCommitJournal: engine.getCommitJournal, + clearCommitJournal: engine.clearCommitJournal, + + // --- Constants --- + + MEMORY_SCHEMA_VERSION: engine.MEMORY_SCHEMA_VERSION, + MEMORY_RELATION_TYPES: engine.MEMORY_RELATION_TYPES, + }); +} diff --git a/src/stats-server.mjs b/src/stats-server.mjs index 5d7221a..95a0ab5 100644 --- a/src/stats-server.mjs +++ b/src/stats-server.mjs @@ -13,17 +13,28 @@ export function isLoopbackAddress(addr) { } /** - * Create an HTTP server that serves memory statistics on GET requests. - * Access is restricted to loopback (localhost) connections only. + * Create an HTTP server that serves memory statistics, telemetry metrics, + * and a health check endpoint. Access is restricted to loopback connections only. * - * @param {{ loadMemories: () => Promise, now?: () => number }} options + * Endpoints: + * GET / — existing memory stats (unchanged) + * GET /metrics — telemetry metrics snapshot (new) + * GET /health — liveness probe (new) + * + * @param {{ loadMemories: () => Promise, now?: () => number, collector?: Object, logger?: Function }} options * @returns {import("node:http").Server} */ -export function createStatsServer({ loadMemories, now = Date.now } = {}) { +export function createStatsServer({ + loadMemories, + now = Date.now, + collector = null, +} = {}) { if (typeof loadMemories !== "function") { throw new TypeError("createStatsServer requires loadMemories to be a function"); } + const serverStartTime = Date.now(); + const server = createServer(async (req, res) => { const remoteAddr = req.socket.remoteAddress; @@ -33,6 +44,49 @@ export function createStatsServer({ loadMemories, now = Date.now } = {}) { return; } + const pathname = new URL(req.url, `http://${req.headers.host || "localhost"}`).pathname; + + // Health endpoint — lightweight liveness probe + if (pathname === "/health") { + try { + const memories = await loadMemories(); + res.writeHead(200, { "Content-Type": "application/json" }); + res.end(JSON.stringify({ + status: "ok", + uptime_ms: Date.now() - serverStartTime, + memory_count: Array.isArray(memories) ? memories.length : 0, + })); + } catch (err) { + res.writeHead(503, { "Content-Type": "application/json" }); + res.end(JSON.stringify({ status: "error", message: err.message })); + } + return; + } + + // Metrics endpoint — telemetry collector snapshot + if (pathname === "/metrics") { + if (!collector) { + res.writeHead(200, { "Content-Type": "application/json" }); + res.end(JSON.stringify({ + message: "No telemetry collector configured", + metrics: {}, + })); + return; + } + try { + const metrics = typeof collector.getMetrics === "function" + ? collector.getMetrics() + : {}; + res.writeHead(200, { "Content-Type": "application/json" }); + res.end(JSON.stringify(metrics)); + } catch (err) { + res.writeHead(500, { "Content-Type": "application/json" }); + res.end(JSON.stringify({ error: "internal_error", message: err.message })); + } + return; + } + + // Default endpoint — existing memory stats behavior try { const memories = await loadMemories(); const stats = computeMemoryStats(memories, { now: now() }); diff --git a/src/telemetry-metrics.mjs b/src/telemetry-metrics.mjs new file mode 100644 index 0000000..c9358da --- /dev/null +++ b/src/telemetry-metrics.mjs @@ -0,0 +1,340 @@ +/** + * Lightweight in-process metrics aggregation for Memact Memory. + * + * Provides counters, histograms, and gauges — all stored in-memory with + * no external dependencies. Designed for integration with the telemetry + * collector and the stats server /metrics endpoint. + * + * @module telemetry-metrics + */ + +/** + * Creates a monotonic counter that can only be incremented. + * @param {string} name - Counter name + * @param {string} [description=""] - Human-readable description + * @returns {Object} Counter instance + */ +export function createCounter(name, description = "") { + let value = 0; + const labels = new Map(); + + return Object.freeze({ + name, + description, + kind: "counter", + + /** + * Increments the counter by the given amount. + * @param {number} [amount=1] + * @param {Object} [labelSet={}] - Optional dimension labels + */ + inc(amount = 1, labelSet = {}) { + const delta = Math.max(0, Number(amount) || 0); + value += delta; + const key = labelKey(labelSet); + if (key) { + labels.set(key, (labels.get(key) || 0) + delta); + } + }, + + /** @returns {number} Current counter value */ + get value() { + return value; + }, + + /** + * Returns value for a specific label set. + * @param {Object} labelSet + * @returns {number} + */ + valueOf(labelSet) { + return labels.get(labelKey(labelSet)) || 0; + }, + + /** Resets the counter to zero. */ + reset() { + value = 0; + labels.clear(); + }, + + /** @returns {Object} Serializable snapshot */ + snapshot() { + return { + name, + kind: "counter", + value, + labels: Object.fromEntries(labels), + }; + }, + }); +} + +/** + * Creates a gauge that can be set to arbitrary values. + * @param {string} name - Gauge name + * @param {string} [description=""] - Human-readable description + * @returns {Object} Gauge instance + */ +export function createGauge(name, description = "") { + let value = 0; + + return Object.freeze({ + name, + description, + kind: "gauge", + + /** + * Sets the gauge to an absolute value. + * @param {number} newValue + */ + set(newValue) { + value = Number(newValue) || 0; + }, + + /** Increments the gauge by amount. */ + inc(amount = 1) { + value += Number(amount) || 0; + }, + + /** Decrements the gauge by amount. */ + dec(amount = 1) { + value -= Number(amount) || 0; + }, + + /** @returns {number} Current gauge value */ + get value() { + return value; + }, + + /** Resets the gauge to zero. */ + reset() { + value = 0; + }, + + /** @returns {Object} Serializable snapshot */ + snapshot() { + return { name, kind: "gauge", value }; + }, + }); +} + +/** + * Creates a histogram for recording value distributions. + * Tracks count, sum, min, max, and configurable percentiles. + * @param {string} name - Histogram name + * @param {Object} [options={}] + * @param {string} [options.description=""] + * @param {number[]} [options.percentiles=[0.5, 0.9, 0.95, 0.99]] - Percentiles to compute + * @param {number} [options.maxSamples=1000] - Max samples to retain for percentile calculation + * @returns {Object} Histogram instance + */ +export function createHistogram(name, options = {}) { + const description = options.description || ""; + const percentiles = options.percentiles || [0.5, 0.9, 0.95, 0.99]; + const maxSamples = options.maxSamples || 1000; + + let count = 0; + let sum = 0; + let min = Infinity; + let max = -Infinity; + let samples = []; + + return Object.freeze({ + name, + description, + kind: "histogram", + + /** + * Records a single observation. + * @param {number} value + */ + observe(value) { + const v = Number(value) || 0; + count += 1; + sum += v; + if (v < min) min = v; + if (v > max) max = v; + samples.push(v); + if (samples.length > maxSamples) { + samples = samples.slice(-maxSamples); + } + }, + + /** @returns {number} Total observation count */ + get count() { + return count; + }, + + /** @returns {number} Sum of all observations */ + get sum() { + return sum; + }, + + /** @returns {number} Minimum observed value */ + get min() { + return count > 0 ? min : 0; + }, + + /** @returns {number} Maximum observed value */ + get max() { + return count > 0 ? max : 0; + }, + + /** @returns {number} Average of all observations */ + get avg() { + return count > 0 ? Number((sum / count).toFixed(4)) : 0; + }, + + /** + * Computes a specific percentile from the retained samples. + * @param {number} p - Percentile as a fraction (e.g. 0.95) + * @returns {number} + */ + percentile(p) { + if (samples.length === 0) return 0; + const sorted = [...samples].sort((a, b) => a - b); + const index = Math.ceil(p * sorted.length) - 1; + return sorted[Math.max(0, index)]; + }, + + /** Resets all histogram state. */ + reset() { + count = 0; + sum = 0; + min = Infinity; + max = -Infinity; + samples = []; + }, + + /** @returns {Object} Serializable snapshot */ + snapshot() { + const pctResults = {}; + for (const p of percentiles) { + const label = `p${Math.round(p * 100)}`; + pctResults[label] = samples.length > 0 ? this.percentile(p) : 0; + } + return { + name, + kind: "histogram", + count, + sum: Number(sum.toFixed(4)), + min: count > 0 ? min : 0, + max: count > 0 ? max : 0, + avg: this.avg, + percentiles: pctResults, + }; + }, + }); +} + +/** + * Creates a complete metrics registry pre-populated with the standard + * Memact Memory metric instruments. + * + * @returns {Object} Registry with named counters, gauges, histograms, and snapshot/reset methods + * + * @example + * const metrics = createMetricsRegistry(); + * metrics.operationTotal.inc(1, { type: "memory.created" }); + * metrics.operationDuration.observe(12.5); + * const snap = metrics.snapshot(); + */ +export function createMetricsRegistry() { + const operationTotal = createCounter( + "operation_total", + "Total operations by type" + ); + const errorTotal = createCounter( + "error_total", + "Total errors by operation type" + ); + const memoryCount = createCounter( + "memory_count", + "Memories processed by type and state" + ); + + const operationDuration = createHistogram("operation_duration_ms", { + description: "Operation duration in milliseconds", + }); + const retrievalScoreDistribution = createHistogram( + "retrieval_score_distribution", + { + description: "Distribution of retrieval scores", + } + ); + const memoryStrengthDistribution = createHistogram( + "memory_strength_distribution", + { + description: "Distribution of memory strength values", + } + ); + + const activeMemoryCount = createGauge( + "active_memory_count", + "Current count of active memories" + ); + const relationCount = createGauge( + "relation_count", + "Current count of memory relations" + ); + const cacheSize = createGauge( + "cache_size", + "Current query cache size" + ); + const databaseSizeBytes = createGauge( + "database_size_bytes", + "Serialized database size in bytes" + ); + + const instruments = { + operationTotal, + errorTotal, + memoryCount, + operationDuration, + retrievalScoreDistribution, + memoryStrengthDistribution, + activeMemoryCount, + relationCount, + cacheSize, + databaseSizeBytes, + }; + + return Object.freeze({ + ...instruments, + + /** + * Returns a serializable snapshot of all instruments. + * @returns {Object} + */ + snapshot() { + const result = {}; + for (const [key, instrument] of Object.entries(instruments)) { + result[key] = instrument.snapshot(); + } + return result; + }, + + /** + * Resets all instruments. Primarily for test isolation. + */ + reset() { + for (const instrument of Object.values(instruments)) { + instrument.reset(); + } + }, + }); +} + +/** + * Produces a stable string key from a label set object for Map lookups. + * @param {Object} labelSet + * @returns {string} + */ +function labelKey(labelSet = {}) { + const entries = Object.entries(labelSet).sort(([a], [b]) => + a.localeCompare(b) + ); + return entries.length > 0 + ? entries.map(([k, v]) => `${k}=${v}`).join(",") + : ""; +} diff --git a/src/telemetry.mjs b/src/telemetry.mjs new file mode 100644 index 0000000..2fb3e95 --- /dev/null +++ b/src/telemetry.mjs @@ -0,0 +1,297 @@ +import { randomBytes } from "node:crypto"; + +/** + * Core telemetry event bus for Memact Memory. + * + * Provides a lightweight pub/sub collector with span timing, typed events, + * wildcard subscribers, buffered flush, and pluggable sink support. + * No external dependencies — designed for zero-overhead opt-in instrumentation. + * + * @module telemetry + */ + +/** + * Canonical telemetry event types emitted by instrumented engine operations. + */ +export const TELEMETRY_EVENTS = Object.freeze({ + MEMORY_CREATED: "memory.created", + MEMORY_UPDATED: "memory.updated", + MEMORY_DELETED: "memory.deleted", + MEMORY_RETRIEVED: "memory.retrieved", + MEMORY_FORGOTTEN: "memory.forgotten", + MEMORY_SUPERSEDED: "memory.superseded", + MEMORY_REINFORCED: "memory.reinforced", + MEMORY_WEAKENED: "memory.weakened", + SCHEMA_ASSIMILATED: "schema.assimilated", + SCHEMA_ACCOMMODATED: "schema.accommodated", + SCHEMA_REMEMBERED: "schema.remembered", + INTENT_REMEMBERED: "intent.remembered", + RELATION_ADDED: "relation.added", + STORE_BUILT: "store.built", + STORE_REINDEXED: "store.reindexed", + QUERY_EXECUTED: "query.executed", + RAG_BUILT: "rag.built", + PACKET_REMEMBERED: "packet.remembered", + CORRECTION_CREATED: "correction.created", + DATABASE_SIZE: "database.size", +}); + +const WILDCARD = "*"; + +/** + * Generates a W3C Trace Context compatible trace ID (32-character hex string). + * @returns {string} + */ +function generateTraceId() { + return randomBytes(16).toString("hex"); +} + +/** + * Creates a new TelemetryEvent object. + * @param {string} type - One of TELEMETRY_EVENTS + * @param {Object} [payload={}] - Event-specific data + * @param {Object} [context={}] - Additional context (trace_id, duration_ms) + * @returns {Object} A frozen telemetry event + */ +function createEvent(type, payload = {}, context = {}) { + return Object.freeze({ + type, + timestamp: new Date().toISOString(), + trace_id: context.trace_id || generateTraceId(), + duration_ms: context.duration_ms ?? null, + memory_id: payload.memory_id || null, + operation: type, + payload, + }); +} + +/** + * Creates a telemetry collector instance with pub/sub event management, + * span timing, buffered flush, and metrics aggregation. + * + * @param {Object} [options={}] + * @param {number} [options.bufferSize=200] - Max events to buffer before auto-flush + * @param {boolean} [options.enabled=true] - Master enable/disable toggle + * @returns {Object} Collector instance + * + * @example + * const collector = createTelemetryCollector(); + * collector.on("memory.created", (event) => console.log(event)); + * collector.on("*", (event) => auditLog(event)); + * + * const span = collector.startSpan("memory.created"); + * // ... do work ... + * span.end({ memory_id: "m_01" }); + */ +export function createTelemetryCollector(options = {}) { + const bufferSize = Number(options.bufferSize ?? 200); + const enabled = options.enabled !== false; + + /** @type {Map>} */ + const subscribers = new Map(); + + /** @type {Array} */ + let eventBuffer = []; + + /** @type {Object} */ + const operationMetrics = {}; + + /** @type {Array} */ + const sinks = []; + + /** + * Registers a subscriber for a specific event type or wildcard. + * @param {string} eventType - Event type or "*" for all events + * @param {Function} handler - Callback receiving the TelemetryEvent + */ + function on(eventType, handler) { + if (typeof handler !== "function") return; + if (!subscribers.has(eventType)) { + subscribers.set(eventType, new Set()); + } + subscribers.get(eventType).add(handler); + } + + /** + * Removes a previously registered subscriber. + * @param {string} eventType + * @param {Function} handler + */ + function off(eventType, handler) { + const handlers = subscribers.get(eventType); + if (handlers) { + handlers.delete(handler); + if (handlers.size === 0) subscribers.delete(eventType); + } + } + + /** + * Emits a telemetry event to all matching subscribers and buffers it. + * @param {string} type - Event type + * @param {Object} [payload={}] - Event data + * @param {Object} [context={}] - Trace context + * @returns {Object} The emitted event + */ + function emit(type, payload = {}, context = {}) { + if (!enabled) return null; + + const event = createEvent(type, payload, context); + + // Update operation metrics + if (!operationMetrics[type]) { + operationMetrics[type] = { count: 0, totalDuration: 0, errors: 0 }; + } + operationMetrics[type].count += 1; + if (event.duration_ms != null) { + operationMetrics[type].totalDuration += event.duration_ms; + } + if (payload.error) { + operationMetrics[type].errors += 1; + } + + // Buffer the event + eventBuffer.push(event); + if (eventBuffer.length > bufferSize) { + eventBuffer = eventBuffer.slice(-bufferSize); + } + + // Notify typed subscribers + const typedHandlers = subscribers.get(type); + if (typedHandlers) { + for (const handler of typedHandlers) { + try { + handler(event); + } catch { + // Subscriber errors must not break the instrumented operation + } + } + } + + // Notify wildcard subscribers + const wildcardHandlers = subscribers.get(WILDCARD); + if (wildcardHandlers) { + for (const handler of wildcardHandlers) { + try { + handler(event); + } catch { + // Subscriber errors must not break the instrumented operation + } + } + } + + return event; + } + + /** + * Starts a timed span for an operation. Call `span.end(payload)` to + * emit the event with duration measurement. + * @param {string} eventType - The event type to emit on end + * @param {Object} [context={}] - Shared trace context + * @returns {{ end: (payload?: Object) => Object }} + */ + function startSpan(eventType, context = {}) { + const traceId = context.trace_id || generateTraceId(); + const startTime = performance.now(); + let ended = false; + + return { + /** + * Ends the span and emits the telemetry event. + * @param {Object} [payload={}] - Event-specific data + * @returns {Object} The emitted event + */ + end(payload = {}) { + if (ended) return null; + ended = true; + const durationMs = Number((performance.now() - startTime).toFixed(2)); + return emit(eventType, payload, { + trace_id: traceId, + duration_ms: durationMs, + }); + }, + }; + } + + /** + * Drains the event buffer to all registered sinks. + * @returns {Array} The flushed events + */ + function flush() { + const events = [...eventBuffer]; + eventBuffer = []; + for (const sink of sinks) { + try { + sink(events); + } catch { + // Sink errors must not propagate + } + } + return events; + } + + /** + * Registers a sink function that receives batches of events on flush. + * @param {Function} sinkFn - Receives an array of events + */ + function addSink(sinkFn) { + if (typeof sinkFn === "function") { + sinks.push(sinkFn); + } + } + + /** + * Returns aggregated metrics for all recorded operations. + * @returns {Object} Metrics keyed by event type + */ + function getMetrics() { + const result = {}; + for (const [type, data] of Object.entries(operationMetrics)) { + result[type] = { + count: data.count, + total_duration_ms: Number(data.totalDuration.toFixed(2)), + avg_duration_ms: + data.count > 0 + ? Number((data.totalDuration / data.count).toFixed(2)) + : 0, + errors: data.errors, + }; + } + return result; + } + + /** + * Returns the current event buffer contents without draining. + * @returns {Array} + */ + function getBuffer() { + return [...eventBuffer]; + } + + /** + * Resets all internal state — metrics, buffer, subscribers, sinks. + * Primarily for test isolation. + */ + function reset() { + eventBuffer = []; + subscribers.clear(); + sinks.length = 0; + for (const key of Object.keys(operationMetrics)) { + delete operationMetrics[key]; + } + } + + return Object.freeze({ + on, + off, + emit, + startSpan, + flush, + addSink, + getMetrics, + getBuffer, + reset, + get enabled() { + return enabled; + }, + }); +} diff --git a/test/audit-trail-test.mjs b/test/audit-trail-test.mjs index cc315cd..481b34b 100644 --- a/test/audit-trail-test.mjs +++ b/test/audit-trail-test.mjs @@ -1,5 +1,6 @@ import assert from "node:assert"; -import { retrieveMemories } from "../src/engine.mjs"; +import { createInstrumentedEngine } from "../src/instrumented-engine.mjs"; +import { createTelemetryCollector, TELEMETRY_EVENTS } from "../src/telemetry.mjs"; console.log("Running Query Audit Trail Compliance Verification..."); @@ -9,15 +10,22 @@ const mockStore = { ] }; -// Execute standard query with context parameters -const results = retrieveMemories("Sample", mockStore, { +// Create a telemetry-instrumented engine and capture events +const collector = createTelemetryCollector(); +const auditEvents = []; +collector.on(TELEMETRY_EVENTS.MEMORY_RETRIEVED, (event) => auditEvents.push(event)); + +const eng = createInstrumentedEngine({ collector }); + +// Execute standard query — audit data is now captured via telemetry events +const results = eng.retrieveMemories("Sample", mockStore, { clientId: "compliance_test_client_44", fieldPath: "user.profile.memories" }); -assert.ok(results.auditTrailLog, "Audit trail token must be appended to retrieval outputs."); -assert.strictEqual(results.auditTrailLog.client_id, "compliance_test_client_44"); -assert.strictEqual(results.auditTrailLog.queried_path, "user.profile.memories"); -assert.ok(typeof results.auditTrailLog.result_count === "number"); +assert.ok(auditEvents.length > 0, "Audit telemetry event must be emitted on retrieval."); +assert.strictEqual(auditEvents[0].type, "memory.retrieved"); +assert.ok(typeof auditEvents[0].payload.result_count === "number"); +assert.ok(auditEvents[0].duration_ms >= 0, "Duration must be recorded."); -console.log("✅ Query audit trail tracking behaves perfectly!"); \ No newline at end of file +console.log("✅ Query audit trail tracking (via telemetry) behaves perfectly!"); \ No newline at end of file diff --git a/test/instrumented-engine.test.mjs b/test/instrumented-engine.test.mjs new file mode 100644 index 0000000..66c9f2f --- /dev/null +++ b/test/instrumented-engine.test.mjs @@ -0,0 +1,275 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { createInstrumentedEngine } from "../src/instrumented-engine.mjs"; +import { createTelemetryCollector, TELEMETRY_EVENTS } from "../src/telemetry.mjs"; +import { createMetricsRegistry } from "../src/telemetry-metrics.mjs"; +import { buildMemoryStore, createMemory, readMemory } from "../src/engine.mjs"; + +/** + * Builds a minimal memory store for testing. + */ +function makeTestStore() { + return buildMemoryStore({ + inference: { + schema_version: "memact.inference.v0", + records: [ + { + id: "rec_01", + packet_id: "pkt_01", + source_label: "Test activity", + meaningful: true, + meaningful_score: 0.82, + sources: [{ url: "https://example.com", title: "Example" }], + canonical_themes: ["testing"], + evidence: { text_excerpt: "Some test evidence" }, + started_at: new Date().toISOString(), + }, + ], + }, + schema: { schemas: [], schema_version: "memact.schema.v0" }, + }); +} + +// Factory + +test("createInstrumentedEngine creates an engine with collector and metrics", () => { + const eng = createInstrumentedEngine(); + assert.ok(eng.collector); + assert.ok(eng.metrics); + assert.strictEqual(typeof eng.createMemory, "function"); + assert.strictEqual(typeof eng.retrieveMemories, "function"); + assert.strictEqual(typeof eng.buildRagContext, "function"); +}); + +test("createInstrumentedEngine accepts injected collector and metrics", () => { + const collector = createTelemetryCollector(); + const metrics = createMetricsRegistry(); + const eng = createInstrumentedEngine({ collector, metrics }); + + assert.strictEqual(eng.collector, collector); + assert.strictEqual(eng.metrics, metrics); +}); + +// Pass-through correctness — results match raw engine + +test("instrumented buildMemoryStore produces same structure as raw engine", () => { + const eng = createInstrumentedEngine(); + const store = eng.buildMemoryStore({ + inference: { + schema_version: "memact.inference.v0", + records: [ + { + id: "rec_01", + packet_id: "pkt_01", + source_label: "Test", + meaningful: true, + meaningful_score: 0.8, + sources: [], + canonical_themes: ["test"], + evidence: {}, + started_at: new Date().toISOString(), + }, + ], + }, + schema: { schemas: [], schema_version: "memact.schema.v0" }, + }); + + assert.ok(Array.isArray(store.memories)); + assert.ok(store.stats); + assert.strictEqual(store.schema_version, "memact.memory.v0"); +}); + +test("instrumented createMemory produces same result as raw engine", () => { + const eng = createInstrumentedEngine(); + const store = makeTestStore(); + + const { memoryStore, memory, action } = eng.createMemory( + { + label: "Test memory", + summary: "A manually created test memory", + strength: 0.7, + }, + store + ); + + assert.ok(memory); + assert.ok(memory.id); + assert.strictEqual(action.accepted, true); + assert.ok(memoryStore.memories.length > store.memories.length); +}); + +test("instrumented retrieveMemories returns same results as raw engine", () => { + const eng = createInstrumentedEngine(); + const store = makeTestStore(); + const results = eng.retrieveMemories("test", store); + + assert.ok(Array.isArray(results)); +}); + +// Telemetry events are emitted + +test("createMemory emits memory.created event", () => { + const collector = createTelemetryCollector(); + const received = []; + collector.on(TELEMETRY_EVENTS.MEMORY_CREATED, (e) => received.push(e)); + + const eng = createInstrumentedEngine({ collector }); + const store = makeTestStore(); + + eng.createMemory({ label: "Observed", strength: 0.6 }, store); + + assert.strictEqual(received.length, 1); + assert.strictEqual(received[0].type, "memory.created"); + assert.ok(received[0].duration_ms >= 0); +}); + +test("retrieveMemories emits memory.retrieved event with result count", () => { + const collector = createTelemetryCollector(); + const received = []; + collector.on(TELEMETRY_EVENTS.MEMORY_RETRIEVED, (e) => received.push(e)); + + const eng = createInstrumentedEngine({ collector }); + const store = makeTestStore(); + + eng.retrieveMemories("test", store); + + assert.strictEqual(received.length, 1); + assert.strictEqual(received[0].type, "memory.retrieved"); + assert.strictEqual(typeof received[0].payload.result_count, "number"); +}); + +test("buildRagContext emits rag.built event", () => { + const collector = createTelemetryCollector(); + const received = []; + collector.on(TELEMETRY_EVENTS.RAG_BUILT, (e) => received.push(e)); + + const eng = createInstrumentedEngine({ collector }); + const store = makeTestStore(); + + eng.buildRagContext("test query", store); + + assert.strictEqual(received.length, 1); + assert.strictEqual(received[0].type, "rag.built"); +}); + +test("deleteMemory emits memory.deleted event", () => { + const collector = createTelemetryCollector(); + const received = []; + collector.on(TELEMETRY_EVENTS.MEMORY_DELETED, (e) => received.push(e)); + + const eng = createInstrumentedEngine({ collector }); + const store = makeTestStore(); + const memoryId = store.memories[0]?.id; + + if (memoryId) { + eng.deleteMemory(memoryId, store, { hard: true }); + assert.strictEqual(received.length, 1); + assert.strictEqual(received[0].type, "memory.deleted"); + } +}); + +test("forgetMemory emits memory.forgotten event", () => { + const collector = createTelemetryCollector(); + const received = []; + collector.on(TELEMETRY_EVENTS.MEMORY_FORGOTTEN, (e) => received.push(e)); + + const eng = createInstrumentedEngine({ collector }); + const store = makeTestStore(); + const memoryId = store.memories[0]?.id; + + if (memoryId) { + eng.forgetMemory(memoryId, store); + assert.strictEqual(received.length, 1); + } +}); + +// Metrics recording + +test("instrumented operations update metrics registry", () => { + const metrics = createMetricsRegistry(); + const eng = createInstrumentedEngine({ metrics }); + const store = makeTestStore(); + + eng.createMemory({ label: "Met1", strength: 0.5 }, store); + eng.createMemory({ label: "Met2", strength: 0.6 }, store); + eng.retrieveMemories("test", store); + + assert.ok(metrics.operationTotal.value >= 3); +}); + +// Duration measurement + +test("span duration is recorded in emitted events", () => { + const collector = createTelemetryCollector(); + const eng = createInstrumentedEngine({ collector }); + const store = makeTestStore(); + + eng.buildMemoryStore({ + inference: { + schema_version: "memact.inference.v0", + records: [ + { + id: "rec_x", + packet_id: "pkt_x", + source_label: "X", + meaningful: true, + meaningful_score: 0.7, + sources: [], + canonical_themes: [], + evidence: {}, + started_at: new Date().toISOString(), + }, + ], + }, + schema: { schemas: [], schema_version: "memact.schema.v0" }, + }); + + const buffer = collector.getBuffer(); + assert.ok(buffer.length > 0); + + const storeEvent = buffer.find((e) => e.type === "store.built"); + assert.ok(storeEvent); + assert.strictEqual(typeof storeEvent.duration_ms, "number"); + assert.ok(storeEvent.duration_ms >= 0); +}); + +// Disabled mode — zero overhead passthrough + +test("disabled instrumented engine passes through without emitting events", () => { + const collector = createTelemetryCollector(); + const received = []; + collector.on("*", (e) => received.push(e)); + + const eng = createInstrumentedEngine({ collector, enabled: false }); + const store = makeTestStore(); + + eng.createMemory({ label: "Silent", strength: 0.5 }, store); + eng.retrieveMemories("test", store); + + // In disabled mode, operations still return correct results + // but no events should have been emitted through the collector + // (the disabled engine uses raw engine functions directly) + assert.strictEqual(received.length, 0); +}); + +// Passthrough functions work unchanged + +test("non-instrumented functions pass through correctly", () => { + const eng = createInstrumentedEngine(); + const store = makeTestStore(); + + // readMemory is a passthrough — should work identically + const memory = eng.readMemory(store.memories[0]?.id, store); + if (store.memories.length > 0) { + assert.ok(memory); + assert.strictEqual(memory.id, store.memories[0].id); + } + + // overlapScore is a passthrough + const score = eng.overlapScore("test", { label: "test", summary: "" }); + assert.strictEqual(typeof score, "number"); + + // Constants are available + assert.strictEqual(eng.MEMORY_SCHEMA_VERSION, "memact.memory.v0"); + assert.ok(eng.MEMORY_RELATION_TYPES); +}); diff --git a/test/telemetry-metrics.test.mjs b/test/telemetry-metrics.test.mjs new file mode 100644 index 0000000..97d415a --- /dev/null +++ b/test/telemetry-metrics.test.mjs @@ -0,0 +1,238 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { + createCounter, + createGauge, + createHistogram, + createMetricsRegistry, +} from "../src/telemetry-metrics.mjs"; + +// Counter + +test("counter starts at zero", () => { + const c = createCounter("test_counter", "A test counter"); + assert.strictEqual(c.value, 0); + assert.strictEqual(c.name, "test_counter"); + assert.strictEqual(c.kind, "counter"); +}); + +test("counter.inc increments by amount", () => { + const c = createCounter("ops"); + c.inc(3); + assert.strictEqual(c.value, 3); + c.inc(); + assert.strictEqual(c.value, 4); +}); + +test("counter.inc ignores negative amounts", () => { + const c = createCounter("ops"); + c.inc(5); + c.inc(-2); + assert.strictEqual(c.value, 5); +}); + +test("counter.inc tracks labeled values", () => { + const c = createCounter("ops"); + c.inc(1, { type: "create" }); + c.inc(1, { type: "create" }); + c.inc(1, { type: "delete" }); + + assert.strictEqual(c.value, 3); + assert.strictEqual(c.valueOf({ type: "create" }), 2); + assert.strictEqual(c.valueOf({ type: "delete" }), 1); + assert.strictEqual(c.valueOf({ type: "update" }), 0); +}); + +test("counter.reset resets to zero", () => { + const c = createCounter("ops"); + c.inc(10); + c.inc(5, { type: "a" }); + c.reset(); + assert.strictEqual(c.value, 0); + assert.strictEqual(c.valueOf({ type: "a" }), 0); +}); + +test("counter.snapshot returns serializable object", () => { + const c = createCounter("ops"); + c.inc(3); + c.inc(2, { region: "us" }); + const snap = c.snapshot(); + + assert.strictEqual(snap.name, "ops"); + assert.strictEqual(snap.kind, "counter"); + assert.strictEqual(snap.value, 5); + assert.strictEqual(typeof snap.labels, "object"); +}); + +// Gauge + +test("gauge starts at zero", () => { + const g = createGauge("memory_count"); + assert.strictEqual(g.value, 0); + assert.strictEqual(g.kind, "gauge"); +}); + +test("gauge.set sets absolute value", () => { + const g = createGauge("memory_count"); + g.set(42); + assert.strictEqual(g.value, 42); + g.set(0); + assert.strictEqual(g.value, 0); +}); + +test("gauge.inc and gauge.dec adjust value", () => { + const g = createGauge("connections"); + g.inc(5); + assert.strictEqual(g.value, 5); + g.dec(2); + assert.strictEqual(g.value, 3); + g.dec(); + assert.strictEqual(g.value, 2); +}); + +test("gauge allows negative values", () => { + const g = createGauge("balance"); + g.dec(10); + assert.strictEqual(g.value, -10); +}); + +test("gauge.reset resets to zero", () => { + const g = createGauge("count"); + g.set(100); + g.reset(); + assert.strictEqual(g.value, 0); +}); + +test("gauge.snapshot returns serializable object", () => { + const g = createGauge("size_bytes", "Database size"); + g.set(1024); + const snap = g.snapshot(); + + assert.strictEqual(snap.name, "size_bytes"); + assert.strictEqual(snap.kind, "gauge"); + assert.strictEqual(snap.value, 1024); +}); + +// Histogram + +test("histogram starts empty", () => { + const h = createHistogram("duration_ms"); + assert.strictEqual(h.count, 0); + assert.strictEqual(h.sum, 0); + assert.strictEqual(h.min, 0); + assert.strictEqual(h.max, 0); + assert.strictEqual(h.avg, 0); +}); + +test("histogram.observe records values correctly", () => { + const h = createHistogram("duration_ms"); + h.observe(10); + h.observe(20); + h.observe(30); + + assert.strictEqual(h.count, 3); + assert.strictEqual(h.sum, 60); + assert.strictEqual(h.min, 10); + assert.strictEqual(h.max, 30); + assert.strictEqual(h.avg, 20); +}); + +test("histogram.percentile calculates correct values", () => { + const h = createHistogram("scores"); + // Add 100 observations: 1, 2, 3, ..., 100 + for (let i = 1; i <= 100; i++) { + h.observe(i); + } + + assert.strictEqual(h.percentile(0.5), 50); + assert.strictEqual(h.percentile(0.9), 90); + assert.strictEqual(h.percentile(0.99), 99); + assert.strictEqual(h.percentile(1.0), 100); +}); + +test("histogram.percentile returns 0 for empty histogram", () => { + const h = createHistogram("empty"); + assert.strictEqual(h.percentile(0.5), 0); +}); + +test("histogram respects maxSamples", () => { + const h = createHistogram("constrained", { maxSamples: 5 }); + for (let i = 0; i < 20; i++) { + h.observe(i); + } + // count tracks all observations, but percentile uses only retained samples + assert.strictEqual(h.count, 20); + // p50 of last 5 values [15,16,17,18,19] should be 17 + assert.strictEqual(h.percentile(0.5), 17); +}); + +test("histogram.reset clears all state", () => { + const h = createHistogram("resettable"); + h.observe(10); + h.observe(20); + h.reset(); + + assert.strictEqual(h.count, 0); + assert.strictEqual(h.sum, 0); + assert.strictEqual(h.min, 0); + assert.strictEqual(h.max, 0); +}); + +test("histogram.snapshot returns serializable object with percentiles", () => { + const h = createHistogram("latency", { + percentiles: [0.5, 0.95], + }); + h.observe(10); + h.observe(20); + h.observe(30); + + const snap = h.snapshot(); + assert.strictEqual(snap.name, "latency"); + assert.strictEqual(snap.kind, "histogram"); + assert.strictEqual(snap.count, 3); + assert.ok("p50" in snap.percentiles); + assert.ok("p95" in snap.percentiles); +}); + +// MetricsRegistry + +test("createMetricsRegistry has all standard instruments", () => { + const reg = createMetricsRegistry(); + + assert.strictEqual(reg.operationTotal.kind, "counter"); + assert.strictEqual(reg.errorTotal.kind, "counter"); + assert.strictEqual(reg.memoryCount.kind, "counter"); + assert.strictEqual(reg.operationDuration.kind, "histogram"); + assert.strictEqual(reg.retrievalScoreDistribution.kind, "histogram"); + assert.strictEqual(reg.memoryStrengthDistribution.kind, "histogram"); + assert.strictEqual(reg.activeMemoryCount.kind, "gauge"); + assert.strictEqual(reg.relationCount.kind, "gauge"); + assert.strictEqual(reg.cacheSize.kind, "gauge"); + assert.strictEqual(reg.databaseSizeBytes.kind, "gauge"); +}); + +test("registry.snapshot returns all instrument snapshots", () => { + const reg = createMetricsRegistry(); + reg.operationTotal.inc(5); + reg.activeMemoryCount.set(42); + reg.operationDuration.observe(12.5); + + const snap = reg.snapshot(); + assert.strictEqual(snap.operationTotal.value, 5); + assert.strictEqual(snap.activeMemoryCount.value, 42); + assert.strictEqual(snap.operationDuration.count, 1); +}); + +test("registry.reset clears all instruments", () => { + const reg = createMetricsRegistry(); + reg.operationTotal.inc(10); + reg.activeMemoryCount.set(50); + reg.operationDuration.observe(100); + + reg.reset(); + + const snap = reg.snapshot(); + assert.strictEqual(snap.operationTotal.value, 0); + assert.strictEqual(snap.activeMemoryCount.value, 0); + assert.strictEqual(snap.operationDuration.count, 0); +}); diff --git a/test/telemetry.test.mjs b/test/telemetry.test.mjs new file mode 100644 index 0000000..50d96a1 --- /dev/null +++ b/test/telemetry.test.mjs @@ -0,0 +1,225 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { + createTelemetryCollector, + TELEMETRY_EVENTS, +} from "../src/telemetry.mjs"; + +// --------------------------------------------------------------------------- +// Collector — creation and event emission +// --------------------------------------------------------------------------- + +test("createTelemetryCollector creates a functional collector", () => { + const collector = createTelemetryCollector(); + assert.strictEqual(collector.enabled, true); + assert.deepStrictEqual(collector.getBuffer(), []); + assert.deepStrictEqual(collector.getMetrics(), {}); +}); + +test("collector.emit creates a well-shaped event", () => { + const collector = createTelemetryCollector(); + const event = collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, { + memory_id: "m_01", + }); + + assert.strictEqual(event.type, "memory.created"); + assert.strictEqual(event.operation, "memory.created"); + assert.strictEqual(event.payload.memory_id, "m_01"); + assert.ok(event.timestamp); + assert.match(event.trace_id, /^[0-9a-f]{32}$/); +}); + +test("collector buffers emitted events", () => { + const collector = createTelemetryCollector(); + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, { memory_id: "m_01" }); + collector.emit(TELEMETRY_EVENTS.MEMORY_UPDATED, { memory_id: "m_02" }); + + const buffer = collector.getBuffer(); + assert.strictEqual(buffer.length, 2); + assert.strictEqual(buffer[0].type, "memory.created"); + assert.strictEqual(buffer[1].type, "memory.updated"); +}); + + +// Collector — subscribers + + + +test("collector.on registers typed subscribers that receive events", () => { + const collector = createTelemetryCollector(); + const received = []; + collector.on(TELEMETRY_EVENTS.MEMORY_CREATED, (e) => received.push(e)); + + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + collector.emit(TELEMETRY_EVENTS.MEMORY_UPDATED, {}); + + assert.strictEqual(received.length, 1); + assert.strictEqual(received[0].type, "memory.created"); +}); + +test("collector.on with wildcard receives all events", () => { + const collector = createTelemetryCollector(); + const received = []; + collector.on("*", (e) => received.push(e)); + + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + collector.emit(TELEMETRY_EVENTS.MEMORY_DELETED, {}); + collector.emit(TELEMETRY_EVENTS.RAG_BUILT, {}); + + assert.strictEqual(received.length, 3); +}); + +test("collector.off removes a subscriber", () => { + const collector = createTelemetryCollector(); + const received = []; + const handler = (e) => received.push(e); + + collector.on(TELEMETRY_EVENTS.MEMORY_CREATED, handler); + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + assert.strictEqual(received.length, 1); + + collector.off(TELEMETRY_EVENTS.MEMORY_CREATED, handler); + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + assert.strictEqual(received.length, 1); +}); + +test("subscriber errors do not break event emission", () => { + const collector = createTelemetryCollector(); + const received = []; + + collector.on(TELEMETRY_EVENTS.MEMORY_CREATED, () => { + throw new Error("subscriber crash"); + }); + collector.on(TELEMETRY_EVENTS.MEMORY_CREATED, (e) => received.push(e)); + + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + assert.strictEqual(received.length, 1); +}); + +// --------------------------------------------------------------------------- +// Collector — span timing +// --------------------------------------------------------------------------- + +test("collector.startSpan measures duration on end", async () => { + const collector = createTelemetryCollector(); + const span = collector.startSpan(TELEMETRY_EVENTS.QUERY_EXECUTED); + + // Simulate some work + await new Promise((resolve) => setTimeout(resolve, 10)); + + const event = span.end({ result_count: 5 }); + assert.strictEqual(event.type, "query.executed"); + assert.ok(event.duration_ms >= 0); + assert.strictEqual(event.payload.result_count, 5); +}); + +test("span.end can only be called once", async () => { + const collector = createTelemetryCollector(); + const span = collector.startSpan(TELEMETRY_EVENTS.MEMORY_CREATED); + + const firstResult = span.end({}); + const secondResult = span.end({}); + + assert.ok(firstResult !== null); + assert.strictEqual(secondResult, null); +}); + +test("span shares trace_id from context", () => { + const collector = createTelemetryCollector(); + const span = collector.startSpan(TELEMETRY_EVENTS.MEMORY_CREATED, { + trace_id: "trace:custom-123", + }); + const event = span.end({}); + assert.strictEqual(event.trace_id, "trace:custom-123"); +}); + +// --------------------------------------------------------------------------- +// Collector — flush and sinks +// --------------------------------------------------------------------------- + +test("collector.flush drains buffer and invokes sinks", () => { + const collector = createTelemetryCollector(); + const sinkBatches = []; + collector.addSink((events) => sinkBatches.push(events)); + + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + collector.emit(TELEMETRY_EVENTS.MEMORY_UPDATED, {}); + + const flushed = collector.flush(); + assert.strictEqual(flushed.length, 2); + assert.strictEqual(sinkBatches.length, 1); + assert.strictEqual(sinkBatches[0].length, 2); + assert.deepStrictEqual(collector.getBuffer(), []); +}); + +test("sink errors do not propagate", () => { + const collector = createTelemetryCollector(); + collector.addSink(() => { + throw new Error("sink crash"); + }); + + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + assert.doesNotThrow(() => collector.flush()); +}); + +// --------------------------------------------------------------------------- +// Collector — metrics +// --------------------------------------------------------------------------- + +test("collector.getMetrics returns operation counts and durations", () => { + const collector = createTelemetryCollector(); + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + collector.emit(TELEMETRY_EVENTS.MEMORY_DELETED, {}); + + const metrics = collector.getMetrics(); + assert.strictEqual(metrics["memory.created"].count, 2); + assert.strictEqual(metrics["memory.deleted"].count, 1); +}); + +test("collector.getMetrics tracks error count", () => { + const collector = createTelemetryCollector(); + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, { error: "test error" }); + + const metrics = collector.getMetrics(); + assert.strictEqual(metrics["memory.created"].errors, 1); +}); + +// --------------------------------------------------------------------------- +// Collector — reset +// --------------------------------------------------------------------------- + +test("collector.reset clears all state", () => { + const collector = createTelemetryCollector(); + collector.on("*", () => {}); + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + + collector.reset(); + assert.deepStrictEqual(collector.getBuffer(), []); + assert.deepStrictEqual(collector.getMetrics(), {}); +}); + +// --------------------------------------------------------------------------- +// Collector — buffer limit +// --------------------------------------------------------------------------- + +test("collector respects buffer size limit", () => { + const collector = createTelemetryCollector({ bufferSize: 3 }); + + for (let i = 0; i < 10; i++) { + collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, { i }); + } + + assert.strictEqual(collector.getBuffer().length, 3); +}); + +// --------------------------------------------------------------------------- +// Disabled collector +// --------------------------------------------------------------------------- + +test("disabled collector returns null from emit", () => { + const collector = createTelemetryCollector({ enabled: false }); + const event = collector.emit(TELEMETRY_EVENTS.MEMORY_CREATED, {}); + assert.strictEqual(event, null); + assert.strictEqual(collector.getBuffer().length, 0); +}); From 97e96dddb0047d23d2686137751c0f9f339d048c Mon Sep 17 00:00:00 2001 From: Shantanu Date: Sat, 18 Jul 2026 21:26:01 +0530 Subject: [PATCH 2/3] Github copilot comments follow-up completed and tested --- src/engine.mjs | 36 ++++++++++++++++++++++-------------- src/instrumented-engine.mjs | 6 +++--- src/stats-server.mjs | 6 +++--- src/telemetry.mjs | 21 +++++++++++++++++---- test/audit-trail-test.mjs | 2 ++ test/telemetry.test.mjs | 21 +++++---------------- 6 files changed, 52 insertions(+), 40 deletions(-) diff --git a/src/engine.mjs b/src/engine.mjs index 09f859e..34759db 100644 --- a/src/engine.mjs +++ b/src/engine.mjs @@ -169,7 +169,7 @@ function tokenSet(value) { function calculateFuzzyMatchScore(stringA, stringB) { const s1 = (stringA || "").toLowerCase().trim(); const s2 = (stringB || "").toLowerCase().trim(); - + if (s1 === s2) return 1.0; if (!s1 || !s2) return 0.0; @@ -183,7 +183,7 @@ function calculateFuzzyMatchScore(stringA, stringB) { for (let i = 0; i < s1.length; i++) { const start = Math.max(0, i - matchWindow); const end = Math.min(s2.length, i + matchWindow + 1); - + for (let j = start; j < end; j++) { if (!s2Matches[j] && s1[i] === s2[j]) { s1Matches[i] = true; @@ -223,7 +223,7 @@ export function overlapScore(query, memory) { const summaryFuzzy = calculateFuzzyMatchScore(queryString, summaryString); const highestFuzzyScore = Math.max(labelFuzzy, summaryFuzzy); - + // Return fuzzy matching score if it meets a reasonable confidence threshold (e.g., > 0.7) return highestFuzzyScore > 0.7 ? highestFuzzyScore : 0.0; } @@ -428,7 +428,7 @@ function decayMemory(memory, options = {}) { const ageDays = daysSince(memory.last_seen_at || memory.first_seen_at); const decay = Math.min(0.35, ageDays * decayPerDay); const decayedStrength = clamp(Number(memory.strength || 0) - decay); - + // Calculate automated TTL expiration trigger thresholds let expirationReason = ""; let state = memory.state || "active"; @@ -618,7 +618,7 @@ function emptyMemoryStore(previous = {}) { * @param {Object} graph * @returns {Object} */ -export function computeStoreStats(memories = [], graph = { nodes: [] }) { +function computeStoreStats(memories = [], graph = { nodes: [] }) { return { memoryCount: memories.length, activityMemoryCount: memories.filter((memory) => memory.type === "activity_memory").length, @@ -849,7 +849,7 @@ export function retrieveMemories(query, memoryStore, options = {}) { searchScore = (lexical * (1 - alpha)) + (semantic * alpha); } const score = clamp((searchScore * 0.56) + (Number(memory.strength || 0) * 0.34) + (isSchemaMemory(memory) ? 0.1 : 0)); - + const BaseResult = { ...memory, retrieval_score: score, @@ -862,7 +862,7 @@ export function retrieveMemories(query, memoryStore, options = {}) { if (auditContext) { const textToAudit = `${query} ${memory.label} ${memory.summary}`; const audit = auditContextLeakage(auditContext, textToAudit); - + BaseResult.audit = { leaked: audit.leaked, violations: audit.violations @@ -877,6 +877,13 @@ export function retrieveMemories(query, memoryStore, options = {}) { .sort((left, right) => right.retrieval_score - left.retrieval_score || right.strength - left.strength) .slice(0, top); + results.auditTrailLog = { + id: `audit:${Date.now()}:${Math.random().toString(36).slice(2, 8)}`, + client_id: clientId, + queried_path: queriedPath, + result_count: results.length, + timestamp: new Date().toISOString(), + }; return results; } @@ -1462,6 +1469,7 @@ function applyMemoryAction(memoryStore, action, mutate) { return mutate(memory); }); const finalAction = matched ? action : { ...action, accepted: false, reason: "memory not found" }; + const nextGraph = buildMemoryGraph(memories, memoryStore.relations || []); const next = { ...memoryStore, memories, @@ -1469,10 +1477,10 @@ function applyMemoryAction(memoryStore, action, mutate) { intent_memories: memories.filter(isIntentMemory), schema_packets: memories.filter(isSchemaMemory), cognitive_schema_memories: memories.filter(isSchemaMemory), - graph: buildMemoryGraph(memories, memoryStore.relations || []), + graph: nextGraph, relations: memoryStore.relations || [], actions: [...(memoryStore.actions || []), finalAction], - stats: computeStoreStats(memories, buildMemoryGraph(memories, memoryStore.relations || [])), + stats: computeStoreStats(memories, nextGraph), }; return { memoryStore: next, action: finalAction }; } @@ -1496,7 +1504,7 @@ function isIntentMemory(memory) { return memory?.type === "intent_memory"; } // Stores historical genre weights to maintain a running baseline profile -let USER_MEDIA_BASELINE = new Map(); +let USER_MEDIA_BASELINE = new Map(); const BASELINE_DECAY = 0.95; // Keeps baseline adaptable but stable const ANOMALY_THRESHOLD = 0.70; // Sensitivity limit for structural signature shifts @@ -1565,18 +1573,18 @@ export function clearMediaBaseline() { */ export function purgeExpiredRecords(records = []) { if (!Array.isArray(records)) return []; - + const currentTime = Date.now(); - + return records.filter(record => { // Check if the record has an expiration or self-destruct timestamp const expirationTime = record.selfDestructAt || record.expiresAt; - + if (expirationTime) { // If the current time has reached or passed expiration, drop the record (return false) return currentTime < new Date(expirationTime).getTime(); } - + // Keep records that don't have an expiration attribute return true; }); diff --git a/src/instrumented-engine.mjs b/src/instrumented-engine.mjs index a852e6a..ef1affa 100644 --- a/src/instrumented-engine.mjs +++ b/src/instrumented-engine.mjs @@ -30,10 +30,10 @@ function instrumentSync(collector, metrics, eventType, fn, extractPayload) { try { const result = fn(...args); const payload = extractPayload ? extractPayload(args, result) : {}; - span.end(payload); + const event = span.end(payload); metrics.operationTotal.inc(1, { type: eventType }); - if (typeof span.end === "function" && payload.duration_ms) { - metrics.operationDuration.observe(payload.duration_ms); + if (event && event.duration_ms != null) { + metrics.operationDuration.observe(event.duration_ms); } return result; } catch (error) { diff --git a/src/stats-server.mjs b/src/stats-server.mjs index 95a0ab5..6c6780d 100644 --- a/src/stats-server.mjs +++ b/src/stats-server.mjs @@ -21,12 +21,12 @@ export function isLoopbackAddress(addr) { * GET /metrics — telemetry metrics snapshot (new) * GET /health — liveness probe (new) * - * @param {{ loadMemories: () => Promise, now?: () => number, collector?: Object, logger?: Function }} options + * @param {{ loadMemories: () => Promise, now?: () => number, collector?: Object }} options * @returns {import("node:http").Server} */ export function createStatsServer({ loadMemories, - now = Date.now, + now = now(), collector = null, } = {}) { if (typeof loadMemories !== "function") { @@ -53,7 +53,7 @@ export function createStatsServer({ res.writeHead(200, { "Content-Type": "application/json" }); res.end(JSON.stringify({ status: "ok", - uptime_ms: Date.now() - serverStartTime, + uptime_ms: now() - serverStartTime, memory_count: Array.isArray(memories) ? memories.length : 0, })); } catch (err) { diff --git a/src/telemetry.mjs b/src/telemetry.mjs index 2fb3e95..f9cc9d7 100644 --- a/src/telemetry.mjs +++ b/src/telemetry.mjs @@ -46,6 +46,19 @@ function generateTraceId() { return randomBytes(16).toString("hex"); } +/** + * Resolves a valid W3C trace ID from context or generates a new one. + * @param {Object} [context={}] + * @returns {string} + */ +function resolveTraceId(context = {}) { + const traceId = context.trace_id; + if (typeof traceId === "string" && /^[0-9a-f]{32}$/i.test(traceId)) { + return traceId; + } + return generateTraceId(); +} + /** * Creates a new TelemetryEvent object. * @param {string} type - One of TELEMETRY_EVENTS @@ -57,7 +70,7 @@ function createEvent(type, payload = {}, context = {}) { return Object.freeze({ type, timestamp: new Date().toISOString(), - trace_id: context.trace_id || generateTraceId(), + trace_id: resolveTraceId(context), duration_ms: context.duration_ms ?? null, memory_id: payload.memory_id || null, operation: type, @@ -70,7 +83,7 @@ function createEvent(type, payload = {}, context = {}) { * span timing, buffered flush, and metrics aggregation. * * @param {Object} [options={}] - * @param {number} [options.bufferSize=200] - Max events to buffer before auto-flush + * @param {number} [options.bufferSize=200] - Max events to retain in buffer (oldest events are dropped when exceeded * @param {boolean} [options.enabled=true] - Master enable/disable toggle * @returns {Object} Collector instance * @@ -84,7 +97,7 @@ function createEvent(type, payload = {}, context = {}) { * span.end({ memory_id: "m_01" }); */ export function createTelemetryCollector(options = {}) { - const bufferSize = Number(options.bufferSize ?? 200); + const bufferSize = Math.max(0, Number(options.bufferSize ?? 200) || 0); const enabled = options.enabled !== false; /** @type {Map>} */ @@ -190,7 +203,7 @@ export function createTelemetryCollector(options = {}) { * @returns {{ end: (payload?: Object) => Object }} */ function startSpan(eventType, context = {}) { - const traceId = context.trace_id || generateTraceId(); + const traceId = resolveTraceId(context); const startTime = performance.now(); let ended = false; diff --git a/test/audit-trail-test.mjs b/test/audit-trail-test.mjs index 481b34b..736aa23 100644 --- a/test/audit-trail-test.mjs +++ b/test/audit-trail-test.mjs @@ -26,6 +26,8 @@ const results = eng.retrieveMemories("Sample", mockStore, { assert.ok(auditEvents.length > 0, "Audit telemetry event must be emitted on retrieval."); assert.strictEqual(auditEvents[0].type, "memory.retrieved"); assert.ok(typeof auditEvents[0].payload.result_count === "number"); +assert.strictEqual(auditEvents[0].payload.client_id, "compliance_test_client_44"); +assert.strictEqual(auditEvents[0].payload.queried_path, "user.profile.memories"); assert.ok(auditEvents[0].duration_ms >= 0, "Duration must be recorded."); console.log("✅ Query audit trail tracking (via telemetry) behaves perfectly!"); \ No newline at end of file diff --git a/test/telemetry.test.mjs b/test/telemetry.test.mjs index 50d96a1..544e597 100644 --- a/test/telemetry.test.mjs +++ b/test/telemetry.test.mjs @@ -5,9 +5,9 @@ import { TELEMETRY_EVENTS, } from "../src/telemetry.mjs"; -// --------------------------------------------------------------------------- + // Collector — creation and event emission -// --------------------------------------------------------------------------- + test("createTelemetryCollector creates a functional collector", () => { const collector = createTelemetryCollector(); @@ -96,9 +96,7 @@ test("subscriber errors do not break event emission", () => { assert.strictEqual(received.length, 1); }); -// --------------------------------------------------------------------------- // Collector — span timing -// --------------------------------------------------------------------------- test("collector.startSpan measures duration on end", async () => { const collector = createTelemetryCollector(); @@ -126,16 +124,15 @@ test("span.end can only be called once", async () => { test("span shares trace_id from context", () => { const collector = createTelemetryCollector(); + const validTraceId = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; const span = collector.startSpan(TELEMETRY_EVENTS.MEMORY_CREATED, { - trace_id: "trace:custom-123", + trace_id: validTraceId, }); const event = span.end({}); - assert.strictEqual(event.trace_id, "trace:custom-123"); + assert.strictEqual(event.trace_id, validTraceId); }); -// --------------------------------------------------------------------------- // Collector — flush and sinks -// --------------------------------------------------------------------------- test("collector.flush drains buffer and invokes sinks", () => { const collector = createTelemetryCollector(); @@ -162,9 +159,7 @@ test("sink errors do not propagate", () => { assert.doesNotThrow(() => collector.flush()); }); -// --------------------------------------------------------------------------- // Collector — metrics -// --------------------------------------------------------------------------- test("collector.getMetrics returns operation counts and durations", () => { const collector = createTelemetryCollector(); @@ -185,9 +180,7 @@ test("collector.getMetrics tracks error count", () => { assert.strictEqual(metrics["memory.created"].errors, 1); }); -// --------------------------------------------------------------------------- // Collector — reset -// --------------------------------------------------------------------------- test("collector.reset clears all state", () => { const collector = createTelemetryCollector(); @@ -199,9 +192,7 @@ test("collector.reset clears all state", () => { assert.deepStrictEqual(collector.getMetrics(), {}); }); -// --------------------------------------------------------------------------- // Collector — buffer limit -// --------------------------------------------------------------------------- test("collector respects buffer size limit", () => { const collector = createTelemetryCollector({ bufferSize: 3 }); @@ -213,9 +204,7 @@ test("collector respects buffer size limit", () => { assert.strictEqual(collector.getBuffer().length, 3); }); -// --------------------------------------------------------------------------- // Disabled collector -// --------------------------------------------------------------------------- test("disabled collector returns null from emit", () => { const collector = createTelemetryCollector({ enabled: false }); From 22ee941b268b414641d9f514808393326de687d4 Mon Sep 17 00:00:00 2001 From: Shantanu Holey <159703391+shantanushok@users.noreply.github.com> Date: Tue, 21 Jul 2026 19:47:55 +0530 Subject: [PATCH 3/3] Implement test for invalid trace_id in telemetry Added test for invalid trace_id handling in span. --- test/telemetry.test.mjs | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/test/telemetry.test.mjs b/test/telemetry.test.mjs index 4ec14e1..7f95b37 100644 --- a/test/telemetry.test.mjs +++ b/test/telemetry.test.mjs @@ -134,6 +134,19 @@ test("span shares trace_id from context", () => { assert.strictEqual(event.trace_id, validTraceId); }); +//checking for invalid trace_id + +test("span rejects invalid trace_id and generates a new W3C-compliant one", () => { + const collector = createTelemetryCollector(); + const invalidTraceId = "trace:custom-123"; + const span = collector.startSpan(TELEMETRY_EVENTS.MEMORY_CREATED, { + trace_id: invalidTraceId, + }); + const event = span.end({}); + assert.notStrictEqual(event.trace_id, invalidTraceId); + assert.match(event.trace_id, /^[0-9a-f]{32}$/i); +}); + // Collector — flush and sinks test("collector.flush drains buffer and invokes sinks", () => {