From 65765bb1c394711e85f089715bcd92667733ceaa Mon Sep 17 00:00:00 2001 From: Thibault Le Ouay Ducasse Date: Thu, 8 Oct 2026 15:04:49 +0200 Subject: [PATCH] server: more analytics --- apps/server/src/libs/cache-keys.ts | 2 + apps/server/src/libs/cli-telemetry.test.ts | 71 +++++++++++ apps/server/src/libs/cli-telemetry.ts | 83 +++++++++++++ apps/server/src/libs/middlewares/track.ts | 39 +++++++ .../interceptors/__tests__/tracking.test.ts | 84 +++++++++++-- .../src/routes/rpc/interceptors/tracking.ts | 110 ++++++++++++------ apps/server/src/routes/rpc/router.ts | 2 +- apps/server/src/routes/v1/index.ts | 7 +- packages/analytics/src/events.ts | 8 ++ 9 files changed, 358 insertions(+), 48 deletions(-) create mode 100644 apps/server/src/libs/cli-telemetry.test.ts create mode 100644 apps/server/src/libs/cli-telemetry.ts diff --git a/apps/server/src/libs/cache-keys.ts b/apps/server/src/libs/cache-keys.ts index 27086569..d6a0ca24 100644 --- a/apps/server/src/libs/cache-keys.ts +++ b/apps/server/src/libs/cache-keys.ts @@ -3,4 +3,6 @@ export const cacheKeys = { pageStatus: (slug: string) => `status:page:${slug}`, monitorDailyStats: (id: string | number) => `stats:monitor:${id}:daily`, + cliInvocation: (workspaceId: number, invocation: string) => + `cli:invocation:${workspaceId}:${invocation}`, }; diff --git a/apps/server/src/libs/cli-telemetry.test.ts b/apps/server/src/libs/cli-telemetry.test.ts new file mode 100644 index 00000000..f832dd81 --- /dev/null +++ b/apps/server/src/libs/cli-telemetry.test.ts @@ -0,0 +1,71 @@ +import { Events } from "@openstatus/analytics"; +import { expect } from "@std/expect"; +import { describe, test } from "@std/testing/bdd"; + +import { claimCliCommandEvent, parseCliHeaders } from "./cli-telemetry"; + +function cliHeaders(overrides: Record = {}) { + return new Headers({ + "user-agent": "openstatus-cli/v1.3.2 (darwin; arm64)", + "x-openstatus-cli-command": "private-locations list", + "x-openstatus-cli-invocation": crypto.randomUUID(), + ...overrides, + }); +} + +describe("parseCliHeaders", () => { + test("reads command, invocation and the user-agent version/platform", () => { + const headers = cliHeaders({ "x-openstatus-cli-invocation": "2452a64f" }); + + expect(parseCliHeaders(headers)).toEqual({ + command: "private-locations list", + invocation: "2452a64f", + version: "1.3.2", + os: "darwin", + arch: "arm64", + }); + }); + + test("returns undefined for requests not sent by the CLI", () => { + expect(parseCliHeaders(new Headers({ "user-agent": "curl/8" }))).toBe( + undefined, + ); + }); + + test("rejects malformed command and invocation values", () => { + expect( + parseCliHeaders(cliHeaders({ "x-openstatus-cli-command": "rm -rf /" })), + ).toBe(undefined); + expect( + parseCliHeaders( + cliHeaders({ "x-openstatus-cli-invocation": "a:b:c:d:e" }), + ), + ).toBe(undefined); + }); + + test("keeps the command when the user agent is not the CLI's", () => { + const cli = parseCliHeaders(cliHeaders({ "user-agent": "custom" })); + + expect(cli?.command).toBe("private-locations list"); + expect(cli?.version).toBe(undefined); + }); +}); + +describe("claimCliCommandEvent", () => { + test("returns the event once per invocation and workspace", async () => { + const cli = parseCliHeaders(cliHeaders()); + if (!cli) throw new Error("expected CLI headers to parse"); + + expect(await claimCliCommandEvent(1, cli)).toEqual({ + ...Events.CliCommand, + command: "private-locations list", + invocation: cli.invocation, + version: "1.3.2", + os: "darwin", + arch: "arm64", + }); + expect(await claimCliCommandEvent(1, cli)).toBe(undefined); + // Keyed per workspace: another workspace can't swallow this run. + expect(await claimCliCommandEvent(2, cli)).not.toBe(undefined); + }); +}); diff --git a/apps/server/src/libs/cli-telemetry.ts b/apps/server/src/libs/cli-telemetry.ts new file mode 100644 index 00000000..a2515c3a --- /dev/null +++ b/apps/server/src/libs/cli-telemetry.ts @@ -0,0 +1,83 @@ +import { getLogger } from "@logtape/logtape"; +import { type EventProps, Events } from "@openstatus/analytics"; + +import { cacheKeys } from "@/libs/cache-keys"; +import { redis } from "@/libs/clients"; + +const logger = getLogger("api-server"); + +export const CLI_COMMAND_HEADER = "x-openstatus-cli-command"; +export const CLI_INVOCATION_HEADER = "x-openstatus-cli-invocation"; + +// `openstatus-cli/v1.3.2 (darwin; arm64)` +const USER_AGENT_RE = + /^openstatus-cli\/v?([\w.+-]{1,32}) \(([^;)]{1,32}); ([^)]{1,32})\)/; +// Full command path with aliases resolved, e.g. `private-locations list`. +const COMMAND_RE = /^[a-z0-9-]+( [a-z0-9-]+){0,4}$/; +// A random per-run id; it becomes part of a Redis key, so keep it tight. +const INVOCATION_RE = /^[A-Za-z0-9-]{8,64}$/; + +// A CLI run is a handful of sequential requests; an hour covers it with room. +const INVOCATION_TTL_SECONDS = 60 * 60; + +export type CliInfo = { + command: string; + invocation: string; + version?: string; + os?: string; + arch?: string; +}; + +/** + * The telemetry the openstatus CLI attaches to every request, or `undefined` + * when the request did not come from the CLI (or carries values we won't + * store). Header values are client-controlled, hence the strict patterns. + */ +export function parseCliHeaders(headers: Headers): CliInfo | undefined { + const command = headers.get(CLI_COMMAND_HEADER); + const invocation = headers.get(CLI_INVOCATION_HEADER); + if (!command || command.length > 100 || !COMMAND_RE.test(command)) return; + if (!invocation || !INVOCATION_RE.test(invocation)) return; + + const ua = USER_AGENT_RE.exec(headers.get("user-agent") ?? ""); + return { + command, + invocation, + ...(ua ? { version: ua[1], os: ua[2], arch: ua[3] } : {}), + }; +} + +/** + * The `cli_command` event for this run, or `undefined` if another request of + * the same run already claimed it. One CLI command can make several requests + * (possibly to different machines); the Redis `NX` claim makes the first one + * win so each run counts once. A Redis failure drops the event rather than + * risk counting a run twice. + */ +export async function claimCliCommandEvent( + workspaceId: number, + cli: CliInfo, +): Promise<(EventProps & Record) | undefined> { + try { + const claimed = await redis.set( + cacheKeys.cliInvocation(workspaceId, cli.invocation), + "1", + { nx: true, ex: INVOCATION_TTL_SECONDS }, + ); + if (!claimed) return; + } catch { + logger.warn("Failed to claim CLI invocation for workspace {workspaceId}", { + workspaceId, + }); + return; + } + + return { + ...Events.CliCommand, + command: cli.command, + invocation: cli.invocation, + ...(cli.version ? { version: cli.version } : {}), + ...(cli.os ? { os: cli.os } : {}), + ...(cli.arch ? { arch: cli.arch } : {}), + }; +} diff --git a/apps/server/src/libs/middlewares/track.ts b/apps/server/src/libs/middlewares/track.ts index 946aea3d..6d54490d 100644 --- a/apps/server/src/libs/middlewares/track.ts +++ b/apps/server/src/libs/middlewares/track.ts @@ -6,6 +6,7 @@ import { } from "@openstatus/analytics"; import type { Context, Next } from "hono"; +import { claimCliCommandEvent, parseCliHeaders } from "@/libs/cli-telemetry"; import type { Variables } from "@/types"; const logger = getLogger("api-server"); @@ -49,3 +50,41 @@ export function trackMiddleware(event: EventProps, eventProps?: string[]) { } }; } + +/** + * Fires `cli_command` for the first request of each openstatus CLI run. V1 is + * deprecated and gets no `api_request` volume tracking, but some CLI commands + * (`whoami`) only ever call it, so their runs would go uncounted without this. + * Mount after `authMiddleware`; counts the run whatever the response status. + */ +export function cliTrackMiddleware() { + return async (c: Context<{ Variables: Variables }, "/*">, next: Next) => { + const cli = parseCliHeaders(c.req.raw.headers); + const workspace = c.get("workspace"); + + if (cli && workspace) { + claimCliCommandEvent(workspace.id, cli) + .then(async (event) => { + if (!event) return; + const analytics = await setupAnalytics({ + userId: `api_${workspace.id}`, + workspaceId: `${workspace.id}`, + workspaceName: workspace.name || workspace.slug, + plan: workspace.plan, + source: "api", + location: c.req.raw.headers.get("x-forwarded-for") ?? undefined, + userAgent: c.req.raw.headers.get("user-agent") ?? undefined, + }); + await analytics.track(event); + }) + .catch(() => { + logger.warn( + "Failed to send CLI analytics event for workspace {workspaceId}", + { workspaceId: workspace.id }, + ); + }); + } + + await next(); + }; +} diff --git a/apps/server/src/routes/rpc/interceptors/__tests__/tracking.test.ts b/apps/server/src/routes/rpc/interceptors/__tests__/tracking.test.ts index bcda2599..6f64e169 100644 --- a/apps/server/src/routes/rpc/interceptors/__tests__/tracking.test.ts +++ b/apps/server/src/routes/rpc/interceptors/__tests__/tracking.test.ts @@ -53,7 +53,7 @@ function createMockRequest( serviceTypeName: string, methodName: string, message: Record = {}, - opts?: { withAuth?: boolean }, + opts?: { withAuth?: boolean; headers?: Record }, ) { const contextValues = new Map(); @@ -71,6 +71,7 @@ function createMockRequest( header: new Headers({ "x-forwarded-for": "1.2.3.4", "user-agent": "test-agent", + ...opts?.headers, }), contextValues: { get: (key: unknown) => contextValues.get(key), @@ -111,7 +112,13 @@ describe("trackingInterceptor", () => { // Flush the .then() chain await Promise.resolve(); - expect(mockTrack).toHaveBeenCalledTimes(1); + expect(mockTrack).toHaveBeenCalledTimes(2); + expect(mockTrack).toHaveBeenCalledWith({ + ...Events.ApiRequest, + service: "openstatus.monitor.v1.MonitorService", + method: "DeleteMonitor", + success: true, + }); expect(mockTrack).toHaveBeenCalledWith({ ...Events.DeleteMonitor, }); @@ -155,22 +162,29 @@ describe("trackingInterceptor", () => { }); }); - test("silently skips unmapped methods", async () => { + test("tracks only api_request for unmapped methods", async () => { const interceptor = trackingInterceptor(); const req = createMockRequest( - "openstatus.health.v1.HealthService", - "Check", + "openstatus.monitor.v1.MonitorService", + "ListMonitors", ); - const mockResponse = { status: "ok" }; + const mockResponse = { monitors: [] }; const next = mockNext(mockResponse); const result = await interceptor(next)(req as never); + await Promise.resolve(); expect(result).toEqual(mockResponse); - expect(mockSetupAnalytics).not.toHaveBeenCalled(); + expect(mockTrack).toHaveBeenCalledTimes(1); + expect(mockTrack).toHaveBeenCalledWith({ + ...Events.ApiRequest, + service: "openstatus.monitor.v1.MonitorService", + method: "ListMonitors", + success: true, + }); }); - test("does not track on error", async () => { + test("tracks failed api_request without the domain event", async () => { const interceptor = trackingInterceptor(); const req = createMockRequest( "openstatus.monitor.v1.MonitorService", @@ -179,7 +193,59 @@ describe("trackingInterceptor", () => { const next = mockNextReject(new Error("not found")); await expect(interceptor(next)(req as never)).rejects.toThrow("not found"); - expect(mockSetupAnalytics).not.toHaveBeenCalled(); + await Promise.resolve(); + + expect(mockSetupAnalytics).toHaveBeenCalledTimes(1); + expect(mockTrack).toHaveBeenCalledTimes(1); + expect(mockTrack).toHaveBeenCalledWith({ + ...Events.ApiRequest, + service: "openstatus.monitor.v1.MonitorService", + method: "DeleteMonitor", + success: false, + }); + }); + + test("tracks one cli_command per CLI invocation", async () => { + const interceptor = trackingInterceptor(); + const invocation = crypto.randomUUID(); + const headers = { + "user-agent": "openstatus-cli/v1.3.2 (darwin; arm64)", + "x-openstatus-cli-command": "monitors list", + "x-openstatus-cli-invocation": invocation, + }; + const flush = () => new Promise((resolve) => setTimeout(resolve, 0)); + + for (let i = 0; i < 2; i++) { + const req = createMockRequest( + "openstatus.monitor.v1.MonitorService", + "ListMonitors", + {}, + { headers }, + ); + await interceptor(mockNext({}))(req as never); + await flush(); + } + + expect(mockTrack).toHaveBeenCalledWith({ + ...Events.ApiRequest, + service: "openstatus.monitor.v1.MonitorService", + method: "ListMonitors", + success: true, + cliCommand: "monitors list", + cliVersion: "1.3.2", + }); + const cliEvents = mockTrack.mock.calls.filter( + ([event]) => event.name === Events.CliCommand.name, + ); + expect(cliEvents).toHaveLength(1); + expect(cliEvents[0][0]).toEqual({ + ...Events.CliCommand, + command: "monitors list", + invocation, + version: "1.3.2", + os: "darwin", + arch: "arm64", + }); }); test("skips tracking when no RPC context", async () => { diff --git a/apps/server/src/routes/rpc/interceptors/tracking.ts b/apps/server/src/routes/rpc/interceptors/tracking.ts index 0511bf96..6f64acf5 100644 --- a/apps/server/src/routes/rpc/interceptors/tracking.ts +++ b/apps/server/src/routes/rpc/interceptors/tracking.ts @@ -7,6 +7,10 @@ import { setupAnalytics, } from "@openstatus/analytics"; +import { + claimCliCommandEvent, + parseCliHeaders, +} from "../../../libs/cli-telemetry"; import { RPC_CONTEXT_KEY } from "./auth"; const logger = getLogger("api-server"); @@ -151,52 +155,84 @@ export const RPC_EVENT_MAP: Record = { /** * Tracking interceptor for ConnectRPC. - * Fires OpenPanel analytics events on successful RPC calls. - * Must be placed after authInterceptor (needs workspace context) - * and validationInterceptor (only track valid requests). + * + * Every authenticated call fires an `api_request` event — reads included, and + * failures too (`success: false`) — so API volume is countable per workspace + * the same way `mcp_request` counts MCP traffic. Successful calls listed in + * `RPC_EVENT_MAP` additionally fire their domain event (e.g. `monitor_created`). + * + * Requests from the openstatus CLI also carry `cliCommand`/`cliVersion` on + * `api_request`, and the first request of each CLI run fires one `cli_command`. + * + * Must be placed after authInterceptor (needs workspace context) and + * validationInterceptor (requests it rejects never reach here, so they are not + * counted). Unauthenticated calls (HealthService) carry no context and are + * skipped. */ export function trackingInterceptor(): Interceptor { return (next) => async (req) => { - const response = await next(req); - - const key = `${req.service.typeName}/${req.method.name}`; - const mapping = RPC_EVENT_MAP[key]; - - if (!mapping) { + let success = false; + try { + const response = await next(req); + success = true; return response; + } finally { + trackRpcCall(req, success); } + }; +} - const rpcCtx = req.contextValues.get(RPC_CONTEXT_KEY); +type TrackedRequest = Parameters[0]>[0]; - if (!rpcCtx) { - return response; - } +function trackRpcCall(req: TrackedRequest, success: boolean) { + const rpcCtx = req.contextValues.get(RPC_CONTEXT_KEY); + if (!rpcCtx) { + return; + } + + const key = `${req.service.typeName}/${req.method.name}`; + const mapping = success ? RPC_EVENT_MAP[key] : undefined; + const cli = parseCliHeaders(req.header); + + const events: (EventProps & Record)[] = [ + { + ...Events.ApiRequest, + service: req.service.typeName, + method: req.method.name, + success, + ...(cli ? { cliCommand: cli.command, cliVersion: cli.version } : {}), + }, + ]; + + if (mapping) { const input = mapping.normalizeInput?.(req.message) ?? req.message; const additionalProps = parseInputToProps(input, mapping.eventProps); - - setupAnalytics({ - userId: `api_${rpcCtx.workspace.id}`, - workspaceId: `${rpcCtx.workspace.id}`, - workspaceName: rpcCtx.workspace.name || rpcCtx.workspace.slug, - plan: rpcCtx.workspace.plan, - source: "api", - location: req.header.get("x-forwarded-for") ?? undefined, - userAgent: req.header.get("user-agent") ?? undefined, + events.push({ ...additionalProps, ...mapping.event }); + } + + const cliEvent = cli + ? claimCliCommandEvent(rpcCtx.workspace.id, cli) + : Promise.resolve(undefined); + + setupAnalytics({ + userId: `api_${rpcCtx.workspace.id}`, + workspaceId: `${rpcCtx.workspace.id}`, + workspaceName: rpcCtx.workspace.name || rpcCtx.workspace.slug, + plan: rpcCtx.workspace.plan, + source: "api", + location: req.header.get("x-forwarded-for") ?? undefined, + userAgent: req.header.get("user-agent") ?? undefined, + }) + .then(async (analytics) => { + const cliCommand = await cliEvent; + if (cliCommand) events.push(cliCommand); + return Promise.all(events.map((event) => analytics.track(event))); }) - .then((analytics) => - analytics.track({ ...additionalProps, ...mapping.event }), - ) - .catch(() => { - logger.warn( - "Failed to send analytics event {event} for workspace {workspaceId}", - { - event: mapping.event.name, - workspaceId: rpcCtx.workspace.id, - }, - ); - }); - - return response; - }; + .catch(() => { + logger.warn( + "Failed to send analytics events for {method} in workspace {workspaceId}", + { method: key, workspaceId: rpcCtx.workspace.id }, + ); + }); } diff --git a/apps/server/src/routes/rpc/router.ts b/apps/server/src/routes/rpc/router.ts index 6233edaa..ba7bf243 100644 --- a/apps/server/src/routes/rpc/router.ts +++ b/apps/server/src/routes/rpc/router.ts @@ -29,7 +29,7 @@ import { * 2. loggingInterceptor - Logs requests/responses with duration * 3. authInterceptor - Validates API key and sets workspace context * 4. validationInterceptor - Validates request messages using protovalidate - * 5. trackingInterceptor - Fires OpenPanel events on success + * 5. trackingInterceptor - Fires OpenPanel `api_request` per call + domain events on success */ export const routes = createConnectRouter({ interceptors: [ diff --git a/apps/server/src/routes/v1/index.ts b/apps/server/src/routes/v1/index.ts index a4afafbb..dc2c35c0 100644 --- a/apps/server/src/routes/v1/index.ts +++ b/apps/server/src/routes/v1/index.ts @@ -6,7 +6,11 @@ import type { RequestIdVariables } from "hono/request-id"; import { env } from "@/env"; import { handleZodError } from "@/libs/errors"; -import { authMiddleware, requireWriteScope } from "@/libs/middlewares"; +import { + authMiddleware, + cliTrackMiddleware, + requireWriteScope, +} from "@/libs/middlewares"; import { checkApi } from "./check"; import { incidentsApi } from "./incidents"; @@ -145,6 +149,7 @@ api.get( * Middlewares */ api.use("/*", authMiddleware); +api.use("/*", cliTrackMiddleware()); // Primary scope enforcement for V1: routes here use inline Drizzle // queries instead of `@openstatus/services`, so the service-level // `requireScope` won't run. After per-route migration to services, diff --git a/packages/analytics/src/events.ts b/packages/analytics/src/events.ts index 07406518..a9cae7bb 100644 --- a/packages/analytics/src/events.ts +++ b/packages/analytics/src/events.ts @@ -316,4 +316,12 @@ export const Events = { name: "mcp_request", channel: "mcp", }, + ApiRequest: { + name: "api_request", + channel: "api", + }, + CliCommand: { + name: "cli_command", + channel: "cli", + }, } as const satisfies Record; -- 2.51.2