Skip to content

Commit 2488a46

Browse files
authored
Merge pull request #3 from syntaxPriest/tool_analytics
Added tools telemetry
2 parents ac99153 + ad553da commit 2488a46

3 files changed

Lines changed: 212 additions & 0 deletions

File tree

0 Bytes
Binary file not shown.

‎mcp/src/analytics.ts‎

Lines changed: 198 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,198 @@
1+
import { createHash, randomUUID } from 'node:crypto'
2+
import * as fs from 'node:fs'
3+
import * as path from 'node:path'
4+
import pkg from '../package.json' with { type: 'json' }
5+
6+
const FLUSH_INTERVAL_MS = 3600_000
7+
const MAX_BATCH_SIZE = 500
8+
const FINAL_FLUSH_TIMEOUT_MS = 5_000
9+
const RECOVERY_DEBOUNCE_MS = 60_000
10+
11+
const PKG_VERSION: string = pkg.version
12+
13+
export interface TelemetryRecord {
14+
v: string
15+
tool: string
16+
ts: number
17+
latencyMs: number
18+
resultSize: number
19+
repoHash: string
20+
hash: string
21+
}
22+
23+
interface BufferedRecord {
24+
v: string
25+
tool: string
26+
ts: number
27+
latencyMs: number
28+
resultSize: number
29+
repoHash: string
30+
hash: string
31+
}
32+
33+
function computeHash(tool: string, latencyMs: number, resultSize: number): string {
34+
return createHash('sha512')
35+
.update(tool + String(latencyMs % 100) + String(resultSize))
36+
.digest('hex')
37+
}
38+
39+
export interface FileBufferedTelemetryOpts {
40+
bufferDir: string
41+
repoHash: string
42+
endpointUrl: string
43+
}
44+
45+
export class FileBufferedTelemetry {
46+
private bufferPath: string
47+
private repoHash: string
48+
private endpointUrl: string
49+
private flushing = false
50+
private intervalHandle: ReturnType<typeof setInterval> | undefined
51+
private closed = false
52+
private lastFlushTs = 0
53+
54+
constructor(opts: FileBufferedTelemetryOpts) {
55+
this.bufferPath = path.join(opts.bufferDir, 'log.jsonl')
56+
this.repoHash = opts.repoHash
57+
this.endpointUrl = opts.endpointUrl
58+
59+
fs.mkdirSync(opts.bufferDir, { recursive: true })
60+
61+
const leftover = this.countLines()
62+
if (leftover > 0) {
63+
const ago = Date.now() - this.getFileMtime()
64+
if (ago >= RECOVERY_DEBOUNCE_MS) {
65+
this.scheduleFlush(100)
66+
}
67+
}
68+
69+
this.intervalHandle = setInterval(() => this.flush(), FLUSH_INTERVAL_MS)
70+
this.intervalHandle.unref()
71+
}
72+
73+
record(tool: string, latencyMs: number, resultSize: number): void {
74+
const rec: BufferedRecord = {
75+
v: PKG_VERSION,
76+
tool,
77+
ts: Date.now(),
78+
latencyMs,
79+
resultSize,
80+
repoHash: this.repoHash,
81+
hash: computeHash(tool, latencyMs, resultSize),
82+
}
83+
try {
84+
fs.appendFileSync(this.bufferPath, JSON.stringify(rec) + '\n', { flag: 'as' })
85+
} catch {
86+
// best-effort
87+
}
88+
}
89+
90+
async close(): Promise<void> {
91+
this.closed = true
92+
if (this.intervalHandle) {
93+
clearInterval(this.intervalHandle)
94+
this.intervalHandle = undefined
95+
}
96+
await this.flushWithTimeout(FINAL_FLUSH_TIMEOUT_MS)
97+
}
98+
99+
private scheduleFlush(delayMs: number): void {
100+
setTimeout(() => { void this.flush() }, delayMs)
101+
}
102+
103+
private countLines(): number {
104+
try {
105+
const content = fs.readFileSync(this.bufferPath, 'utf8')
106+
if (content.length === 0) return 0
107+
return content.split('\n').filter(Boolean).length
108+
} catch {
109+
return 0
110+
}
111+
}
112+
113+
private getFileMtime(): number {
114+
try {
115+
return fs.statSync(this.bufferPath).mtimeMs
116+
} catch {
117+
return 0
118+
}
119+
}
120+
121+
private readRecords(): BufferedRecord[] {
122+
try {
123+
const content = fs.readFileSync(this.bufferPath, 'utf8')
124+
if (!content) return []
125+
const lines = content.split('\n').filter(Boolean)
126+
const records: BufferedRecord[] = []
127+
for (const line of lines) {
128+
try {
129+
records.push(JSON.parse(line) as BufferedRecord)
130+
} catch {
131+
// skip malformed lines
132+
}
133+
}
134+
return records
135+
} catch {
136+
return []
137+
}
138+
}
139+
140+
private async flushWithTimeout(timeoutMs: number): Promise<void> {
141+
const result = this.flush()
142+
const timer = new Promise<void>((_, reject) => {
143+
setTimeout(() => reject(new Error('timeout')), timeoutMs)
144+
})
145+
try {
146+
await Promise.race([result, timer])
147+
} catch {
148+
// timeout — move on
149+
}
150+
}
151+
152+
private async flush(): Promise<void> {
153+
if (this.flushing) return
154+
this.flushing = true
155+
this.lastFlushTs = Date.now()
156+
try {
157+
const records = this.readRecords()
158+
if (records.length === 0) return
159+
160+
const batch = records.slice(0, MAX_BATCH_SIZE)
161+
162+
await this.upload(batch)
163+
164+
if (!this.closed) {
165+
const remaining = records.slice(MAX_BATCH_SIZE)
166+
const content = remaining.map((r) => JSON.stringify(r)).join('\n') + (remaining.length > 0 ? '\n' : '')
167+
fs.writeFileSync(this.bufferPath, content, 'utf8')
168+
}
169+
} catch (err) {
170+
const msg = err instanceof Error ? err.message : String(err)
171+
process.stderr.write(`openvisio telemetry: flush failed (${msg})\n`)
172+
} finally {
173+
this.flushing = false
174+
}
175+
}
176+
177+
private async upload(records: BufferedRecord[]): Promise<void> {
178+
const body = JSON.stringify(records)
179+
const ctrl = new AbortController()
180+
const t = setTimeout(() => ctrl.abort(), 10_000)
181+
try {
182+
const res = await fetch(this.endpointUrl, {
183+
method: 'POST',
184+
headers: { 'Content-Type': 'application/json' },
185+
body,
186+
signal: ctrl.signal,
187+
})
188+
if (!res.ok) {
189+
throw new Error(`HTTP ${res.status}`)
190+
}
191+
if (!this.closed) {
192+
fs.truncateSync(this.bufferPath, 0)
193+
}
194+
} finally {
195+
clearTimeout(t)
196+
}
197+
}
198+
}

‎mcp/src/server.ts‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,12 +5,14 @@
55
// `--spotlight` lights up an open viewer. Local-first, read-only; only the
66
// spotlight binds a local (127.0.0.1) port.
77

8+
import { createHash } from 'node:crypto'
89
import * as fs from 'node:fs'
910
import * as os from 'node:os'
1011
import * as path from 'node:path'
1112
import { McpServer } from '@modelcontextprotocol/sdk/server/mcp.js'
1213
import { StdioServerTransport } from '@modelcontextprotocol/sdk/server/stdio.js'
1314
import { computeCentrality, computeChurn, Indexer, type CodeGraph } from '@openvisio/core'
15+
import { FileBufferedTelemetry } from './analytics.js'
1416
import { SavingsReceipt } from './receipt.js'
1517
import { startSpotlightServer, type SpotlightEvent, type SpotlightServer, type UserRequest } from './spotlight.js'
1618
import { buildTools, type GraphState } from './tools.js'
@@ -88,6 +90,13 @@ export async function serveMcp(opts: ServeOptions): Promise<void> {
8890
// index everything under it — effectively forever. Refuse, with a clear message.
8991
const badRoot = resolvedRoot === path.parse(resolvedRoot).root || resolvedRoot === path.resolve(os.homedir())
9092

93+
const repoHash = createHash('sha256').update(resolvedRoot).digest('hex').slice(0, 12)
94+
const telemetry = new FileBufferedTelemetry({
95+
bufferDir: path.join(os.homedir(), '.local', 'share', 'openvisio', 'telemetry'),
96+
repoHash,
97+
endpointUrl: 'https://k5b3bh9hte.execute-api.us-east-1.amazonaws.com/dev/tool/telemetary',
98+
})
99+
91100
const server = new McpServer(
92101
{ name: 'openvisio', version: '0.1.5' },
93102
{
@@ -148,8 +157,11 @@ export async function serveMcp(opts: ServeOptions): Promise<void> {
148157
{ description: tool.description, inputSchema: tool.inputShape },
149158
async (args: Record<string, unknown>) => {
150159
await ready
160+
const t0 = Date.now()
151161
try {
152162
const result = tool.handler(args)
163+
const latencyMs = Date.now() - t0
164+
telemetry.record(tool.name, latencyMs, result.text.length)
153165
receipt.record(result.text, result.touchedFiles)
154166
if (spotlight && result.touchedFiles.length > 0) {
155167
const { focus, edges } = spotlightPayload(getState().graph, result.touchedFiles)
@@ -158,6 +170,7 @@ export async function serveMcp(opts: ServeOptions): Promise<void> {
158170
return { content: [{ type: 'text' as const, text: result.text }] }
159171
} catch (err) {
160172
const msg = err instanceof Error ? err.message : String(err)
173+
telemetry.record(tool.name, Date.now() - t0, 0)
161174
return { content: [{ type: 'text' as const, text: `Error: ${msg}` }], isError: true }
162175
}
163176
},
@@ -178,6 +191,7 @@ export async function serveMcp(opts: ServeOptions): Promise<void> {
178191
watcher?.close()
179192
spotlight?.close()
180193
indexer.close()
194+
telemetry.close()
181195
if (state) {
182196
const summary = receipt.summary()
183197
if (summary) process.stderr.write(summary + '\n')

0 commit comments

Comments
 (0)