From 4b467e8b2b1ee633f178f220ca14b725f880f4c8 Mon Sep 17 00:00:00 2001 From: Thibault Le Ouay Date: Thu, 8 Oct 2026 17:01:52 +0200 Subject: [PATCH] server: more analytics (#2872) * server: more analytics * update pr * update pr * Update apps/server/src/libs/analytics-identity.ts Co-authored-by: cubic-dev-ai[bot] <191113872+cubic-dev-ai[bot]@users.noreply.github.com> * update pr --------- Co-authored-by: cubic-dev-ai[bot] <191113872+cubic-dev-ai[bot]@users.noreply.github.com> --- apps/server/src/libs/analytics-identity.ts | 24 ++++ apps/server/src/libs/cache-keys.ts | 2 + apps/server/src/libs/cli-telemetry.test.ts | 96 +++++++++++++ apps/server/src/libs/cli-telemetry.ts | 112 +++++++++++++++ .../server/src/libs/middlewares/track.test.ts | 135 ++++++++++++++++++ apps/server/src/libs/middlewares/track.ts | 65 +++++++-- .../src/libs/test/doubles/analytics.mock.ts | 2 + .../interceptors/__tests__/tracking.test.ts | 96 +++++++++++-- .../src/routes/rpc/interceptors/tracking.ts | 108 +++++++++----- apps/server/src/routes/rpc/router.ts | 2 +- apps/server/src/routes/v1/index.ts | 7 +- packages/analytics/src/events.ts | 8 ++ packages/analytics/src/server.ts | 62 +++++++- 13 files changed, 649 insertions(+), 70 deletions(-) create mode 100644 apps/server/src/libs/analytics-identity.ts create mode 100644 apps/server/src/libs/cli-telemetry.test.ts create mode 100644 apps/server/src/libs/cli-telemetry.ts create mode 100644 apps/server/src/libs/middlewares/track.test.ts diff --git a/apps/server/src/libs/analytics-identity.ts b/apps/server/src/libs/analytics-identity.ts new file mode 100644 index 00000000..3f56d311 --- /dev/null +++ b/apps/server/src/libs/analytics-identity.ts @@ -0,0 +1,24 @@ +import type { IdentifyProps } from "@openstatus/analytics"; +import type { Workspace } from "@openstatus/db/src/schema"; + +/** + * The OpenPanel identity for API traffic: one `api_` profile per + * workspace, shared by the V1 REST and ConnectRPC surfaces. + */ +export function apiAnalyticsIdentity( + workspace: Workspace, + headers: Headers, +): IdentifyProps { + return { + userId: `api_${workspace.id}`, + workspaceId: `${workspace.id}`, + workspaceName: workspace.name || workspace.slug, + plan: workspace.plan, + source: "api", + location: + headers.get("fly-client-ip") || + headers.get("x-forwarded-for")?.split(",").at(-1)?.trim() || + undefined, + userAgent: headers.get("user-agent") ?? undefined, + }; +} 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..e072dedb --- /dev/null +++ b/apps/server/src/libs/cli-telemetry.test.ts @@ -0,0 +1,96 @@ +import { Events } from "@openstatus/analytics"; +import { expect } from "@std/expect"; +import { describe, test } from "@std/testing/bdd"; + +import { + claimCliCommandEvent, + parseCliHeaders, + trackCliCommand, +} 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); + }); +}); + +describe("trackCliCommand", () => { + test("releases the claim when the send fails so a later request retries", async () => { + const cli = parseCliHeaders(cliHeaders()); + if (!cli) throw new Error("expected CLI headers to parse"); + const sent: unknown[] = []; + const failing = { track: () => Promise.reject(new Error("down")) }; + const working = { + track: (event: unknown) => { + sent.push(event); + return Promise.resolve(); + }, + }; + + await expect(trackCliCommand(failing, 1, cli)).rejects.toThrow("down"); + await trackCliCommand(working, 1, cli); + await trackCliCommand(working, 1, cli); + + expect(sent).toHaveLength(1); + }); +}); diff --git a/apps/server/src/libs/cli-telemetry.ts b/apps/server/src/libs/cli-telemetry.ts new file mode 100644 index 00000000..4fd277f0 --- /dev/null +++ b/apps/server/src/libs/cli-telemetry.ts @@ -0,0 +1,112 @@ +import { getLogger } from "@logtape/logtape"; +import { type Analytics, 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. Prefer `trackCliCommand`, which also gives the + * claim back when the send fails. + */ +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 } : {}), + }; +} + +/** + * Sends this run's `cli_command` through `analytics` if this request is the + * first of the run to claim it. Claim only once analytics is set up, and give + * the claim back if the send rejects, so a later request of the same run can + * still count it instead of the run being lost for the claim's TTL. + */ +export async function trackCliCommand( + analytics: Analytics, + workspaceId: number, + cli: CliInfo, +): Promise { + const event = await claimCliCommandEvent(workspaceId, cli); + if (!event) return; + try { + await analytics.track(event); + } catch (error) { + await redis + .del(cacheKeys.cliInvocation(workspaceId, cli.invocation)) + .catch(() => { + logger.warn( + "Failed to release CLI invocation claim for workspace {workspaceId}", + { workspaceId }, + ); + }); + throw error; + } +} diff --git a/apps/server/src/libs/middlewares/track.test.ts b/apps/server/src/libs/middlewares/track.test.ts new file mode 100644 index 00000000..bf21cb12 --- /dev/null +++ b/apps/server/src/libs/middlewares/track.test.ts @@ -0,0 +1,135 @@ +import { Events } from "@openstatus/analytics"; +import type { Workspace } from "@openstatus/db/src/schema"; +import { + beforeEach, + describe, + expect, + type mock, + test, +} from "@openstatus/test-utils"; +import { Hono } from "hono"; + +import type { Variables } from "@/types"; + +import { apiTrackMiddleware } from "./track"; + +// @openstatus/analytics is swapped for a double (test.importmap.json) whose +// setupAnalytics/track spies are exposed here on globalThis. +const { track: mockTrack, setupAnalytics: mockSetupAnalytics } = ( + globalThis as Record +).__analyticsSpies as { + track: ReturnType; + setupAnalytics: ReturnType; +}; + +/** Let the fire-and-forget analytics chain settle. */ +const flush = () => new Promise((resolve) => setTimeout(resolve, 0)); + +const TEST_WORKSPACE = { + id: 42, + slug: "test-ws", + plan: "free", + name: "", +} as Workspace; + +/** + * Minimal app with `workspace` pre-set (simulating `authMiddleware`) and + * `apiTrackMiddleware` mounted the way the V1 router mounts it. + */ +function makeApp(opts: { withWorkspace?: boolean } = {}) { + const app = new Hono<{ Variables: Variables }>(); + app.use("*", async (c, next) => { + if (opts.withWorkspace !== false) c.set("workspace", TEST_WORKSPACE); + await next(); + }); + app.use("*", apiTrackMiddleware()); + app.get("/whoami", (c) => c.text("ok")); + app.get("/monitor/:id", (c) => c.text("missing", 404)); + return app; +} + +function cliHeaders(invocation: string) { + return { + "user-agent": "openstatus-cli/v1.3.2 (darwin; arm64)", + "x-openstatus-cli-command": "whoami", + "x-openstatus-cli-invocation": invocation, + }; +} + +function cliCommandCalls() { + return mockTrack.mock.calls.filter( + ([event]) => (event as { name: string }).name === Events.CliCommand.name, + ); +} + +describe("apiTrackMiddleware", () => { + beforeEach(() => { + mockSetupAnalytics.mockClear(); + mockTrack.mockClear(); + }); + + test("fires api_request for GET requests with the route pattern", async () => { + const res = await makeApp().request("/monitor/123", { + headers: { "user-agent": "curl/8" }, + }); + await flush(); + + expect(res.status).toBe(404); + expect(mockTrack).toHaveBeenCalledTimes(1); + expect(mockTrack).toHaveBeenCalledWith({ + ...Events.ApiRequest, + service: "v1", + method: "GET /monitor/:id", + success: false, + }); + }); + + test("skips requests that match no route", async () => { + const res = await makeApp().request("/nope", { + headers: { "user-agent": "curl/8" }, + }); + await flush(); + + expect(res.status).toBe(404); + expect(mockSetupAnalytics).not.toHaveBeenCalled(); + }); + + test("skips requests without a workspace", async () => { + const res = await makeApp({ withWorkspace: false }).request("/whoami", { + headers: cliHeaders(crypto.randomUUID()), + }); + await flush(); + + expect(res.status).toBe(200); + expect(mockSetupAnalytics).not.toHaveBeenCalled(); + expect(mockTrack).not.toHaveBeenCalled(); + }); + + test("tags CLI requests and fires cli_command once per invocation", async () => { + const app = makeApp(); + const invocation = crypto.randomUUID(); + + for (let i = 0; i < 2; i++) { + await app.request("/whoami", { headers: cliHeaders(invocation) }); + await flush(); + } + + expect(mockTrack).toHaveBeenCalledWith({ + ...Events.ApiRequest, + service: "v1", + method: "GET /whoami", + success: true, + cliCommand: "whoami", + cliVersion: "1.3.2", + }); + expect(cliCommandCalls()).toHaveLength(1); + expect(cliCommandCalls()[0][0]).toEqual({ + ...Events.CliCommand, + command: "whoami", + invocation, + version: "1.3.2", + os: "darwin", + arch: "arm64", + }); + }); +}); diff --git a/apps/server/src/libs/middlewares/track.ts b/apps/server/src/libs/middlewares/track.ts index 946aea3d..628f48d4 100644 --- a/apps/server/src/libs/middlewares/track.ts +++ b/apps/server/src/libs/middlewares/track.ts @@ -1,11 +1,15 @@ import { getLogger } from "@logtape/logtape"; import { type EventProps, + Events, parseInputToProps, setupAnalytics, } from "@openstatus/analytics"; import type { Context, Next } from "hono"; +import { matchedRoutes } from "hono/route"; +import { apiAnalyticsIdentity } from "@/libs/analytics-identity"; +import { parseCliHeaders, trackCliCommand } from "@/libs/cli-telemetry"; import type { Variables } from "@/types"; const logger = getLogger("api-server"); @@ -30,15 +34,7 @@ export function trackMiddleware(event: EventProps, eventProps?: string[]) { const additionalProps = parseInputToProps(json, eventProps); const workspace = c.get("workspace"); - 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, - }) + setupAnalytics(apiAnalyticsIdentity(workspace, c.req.raw.headers)) .then((analytics) => analytics.track({ ...additionalProps, ...event })) .catch(() => { logger.warn( @@ -49,3 +45,54 @@ export function trackMiddleware(event: EventProps, eventProps?: string[]) { } }; } + +/** + * Fires `api_request` for every authenticated V1 call — reads included, and + * failures too (`success: false`) — mirroring the RPC tracking interceptor so + * API volume is countable across both surfaces. `method` is the matched route + * pattern (`GET /v1/monitor/:id`), never the raw URL, to keep it low-cardinality. + * + * Requests from the openstatus CLI also carry `cliCommand`/`cliVersion`, and + * the first request of each CLI run fires one `cli_command`. + * + * Mount after `authMiddleware`; requests it rejects carry no workspace and are + * skipped, as are requests that match no route. Per-route domain events stay with `trackMiddleware`. + */ +export function apiTrackMiddleware() { + return async (c: Context<{ Variables: Variables }, "/*">, next: Next) => { + await next(); + + const workspace = c.get("workspace"); + if (!workspace) return; + + // `use()` middlewares register as `ALL`; with no handler route matched + // (a 404 on an unknown path) the only pattern left is `/*`, which says + // nothing about the endpoint, so those aren't counted. + const route = matchedRoutes(c).findLast((r) => r.method !== "ALL"); + if (!route) return; + + const cli = parseCliHeaders(c.req.raw.headers); + const success = c.res.status.toString().startsWith("2") && !c.error; + + setupAnalytics(apiAnalyticsIdentity(workspace, c.req.raw.headers)) + .then((analytics) => + Promise.all([ + analytics.track({ + ...Events.ApiRequest, + service: "v1", + method: `${c.req.method} ${route.path}`, + success, + ...(cli ? { cliCommand: cli.command } : {}), + ...(cli?.version ? { cliVersion: cli.version } : {}), + }), + ...(cli ? [trackCliCommand(analytics, workspace.id, cli)] : []), + ]), + ) + .catch(() => { + logger.warn( + "Failed to send API analytics events for workspace {workspaceId}", + { workspaceId: workspace.id }, + ); + }); + }; +} diff --git a/apps/server/src/libs/test/doubles/analytics.mock.ts b/apps/server/src/libs/test/doubles/analytics.mock.ts index dfe71343..3ac49651 100644 --- a/apps/server/src/libs/test/doubles/analytics.mock.ts +++ b/apps/server/src/libs/test/doubles/analytics.mock.ts @@ -2,8 +2,10 @@ // real Events/parseInputToProps but replaces setupAnalytics with a spy (exposed // on globalThis.__analyticsSpies) so tests assert tracking without OpenPanel. export { + type Analytics, type EventProps, Events, + type IdentifyProps, parseInputToProps, } from "@openstatus/analytics-real"; 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..baa05699 100644 --- a/apps/server/src/routes/rpc/interceptors/__tests__/tracking.test.ts +++ b/apps/server/src/routes/rpc/interceptors/__tests__/tracking.test.ts @@ -25,6 +25,9 @@ const { track: mockTrack, setupAnalytics: mockSetupAnalytics } = ( type NextFn = Parameters>[0]; +/** Let the fire-and-forget analytics chain (several awaits deep) settle. */ +const flush = () => new Promise((resolve) => setTimeout(resolve, 0)); + /** Create a mock `next` that resolves with the given value. */ function mockNext(response: unknown): NextFn { return mock(() => Promise.resolve(response)) as unknown as NextFn; @@ -53,7 +56,7 @@ function createMockRequest( serviceTypeName: string, methodName: string, message: Record = {}, - opts?: { withAuth?: boolean }, + opts?: { withAuth?: boolean; headers?: Record }, ) { const contextValues = new Map(); @@ -71,6 +74,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), @@ -109,9 +113,15 @@ describe("trackingInterceptor", () => { }); // Flush the .then() chain - await Promise.resolve(); + await flush(); - 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, }); @@ -127,7 +137,7 @@ describe("trackingInterceptor", () => { const next = mockNext({}); await interceptor(next)(req as never); - await Promise.resolve(); + await flush(); expect(mockTrack).toHaveBeenCalledWith({ ...Events.CreateMonitor, @@ -146,7 +156,7 @@ describe("trackingInterceptor", () => { const next = mockNext({}); await interceptor(next)(req as never); - await Promise.resolve(); + await flush(); expect(mockTrack).toHaveBeenCalledWith({ ...Events.CreateMonitor, @@ -155,22 +165,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 flush(); 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 +196,58 @@ 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 flush(); + + 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, + }; + + 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 () => { @@ -214,7 +282,7 @@ describe("trackingInterceptor", () => { expect(result).toEqual(mockResponse); // Flush the .catch() chain — should not throw - await Promise.resolve(); + await flush(); expect(mockTrack).not.toHaveBeenCalled(); }); @@ -241,7 +309,7 @@ describe("trackingInterceptor", () => { }); // Flush the .then() chain - await Promise.resolve(); + await flush(); expect(mockTrack).toHaveBeenCalledWith({ ...Events.CreateNotification, diff --git a/apps/server/src/routes/rpc/interceptors/tracking.ts b/apps/server/src/routes/rpc/interceptors/tracking.ts index 0511bf96..4db826c1 100644 --- a/apps/server/src/routes/rpc/interceptors/tracking.ts +++ b/apps/server/src/routes/rpc/interceptors/tracking.ts @@ -7,6 +7,8 @@ import { setupAnalytics, } from "@openstatus/analytics"; +import { apiAnalyticsIdentity } from "../../../libs/analytics-identity"; +import { parseCliHeaders, trackCliCommand } from "../../../libs/cli-telemetry"; import { RPC_CONTEXT_KEY } from "./auth"; const logger = getLogger("api-server"); @@ -151,52 +153,82 @@ 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. V1 REST fires the same event + * from `apiTrackMiddleware`. 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 { + // Tracking must never replace the call's own result or error. + try { + trackRpcCall(req, success); + } catch { + logger.warn("Failed to track RPC call {method}", { + method: `${req.service.typeName}/${req.method.name}`, + }); + } } + }; +} - 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 } : {}), + ...(cli?.version ? { 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, - }) - .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; - }; + events.push({ ...additionalProps, ...mapping.event }); + } + + setupAnalytics(apiAnalyticsIdentity(rpcCtx.workspace, req.header)) + .then((analytics) => + Promise.all([ + ...events.map((event) => analytics.track(event)), + ...(cli ? [trackCliCommand(analytics, rpcCtx.workspace.id, cli)] : []), + ]), + ) + .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..8dfb4bde 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` for authenticated, validated calls + 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..3866c045 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, + apiTrackMiddleware, + requireWriteScope, +} from "@/libs/middlewares"; import { checkApi } from "./check"; import { incidentsApi } from "./incidents"; @@ -145,6 +149,7 @@ api.get( * Middlewares */ api.use("/*", authMiddleware); +api.use("/*", apiTrackMiddleware()); // 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; diff --git a/packages/analytics/src/server.ts b/packages/analytics/src/server.ts index 59811929..fcff971a 100644 --- a/packages/analytics/src/server.ts +++ b/packages/analytics/src/server.ts @@ -47,6 +47,42 @@ export type IdentifyProps = { userAgent?: string; }; +// identify, upsertGroup and setGroup are an HTTP round trip each, and +// per-request callers (every API and MCP call) send the same payload over and +// over. Remember which identities this process sent recently and skip them; +// a changed name or plan is a new key, so it still goes out. +const IDENTITY_TTL_MS = 60 * 60 * 1000; +const MAX_IDENTITIES = 10_000; +const sentIdentities = new Map(); + +function identityKey(props: IdentifyProps) { + return JSON.stringify([ + props.userId, + props.fullName, + props.email, + props.workspaceId, + props.workspaceName, + props.plan, + ]); +} + +function wasIdentitySent(key: string) { + const expiresAt = sentIdentities.get(key); + if (expiresAt === undefined) return false; + if (expiresAt > Date.now()) return true; + sentIdentities.delete(key); + return false; +} + +function markIdentitySent(key: string) { + // Map iterates in insertion order, so the first key is the oldest. + if (sentIdentities.size >= MAX_IDENTITIES) { + const oldest = sentIdentities.keys().next().value; + if (oldest !== undefined) sentIdentities.delete(oldest); + } + sentIdentities.set(key, Date.now() + IDENTITY_TTL_MS); +} + export async function setupAnalytics(props: IdentifyProps) { if (env.NODE_ENV !== "production") { return noop(); @@ -62,6 +98,21 @@ export async function setupAnalytics(props: IdentifyProps) { op.api.addHeader("user-agent", props.userAgent); } + const groupId = props.workspaceId ? `ws_${props.workspaceId}` : undefined; + // profileId and groups go on each event rather than relying on the client + // state identify()/setGroup() leave behind, which the cache below skips. + const track = (opts: EventProps & TrackProperties) => { + const { name, ...rest } = opts; + return op.track(name, { + ...(props.userId ? { profileId: props.userId } : {}), + ...rest, + ...(groupId ? { groups: [groupId] } : {}), + }); + }; + + const key = identityKey(props); + if (wasIdentitySent(key)) return { track }; + if (props.userId) { const [firstName, lastName] = props.fullName?.split(" ") || []; await op.identify({ @@ -72,7 +123,6 @@ export async function setupAnalytics(props: IdentifyProps) { }); } - const groupId = props.workspaceId ? `ws_${props.workspaceId}` : undefined; if (groupId && props.workspaceName) { await op.upsertGroup({ id: groupId, @@ -85,15 +135,13 @@ export async function setupAnalytics(props: IdentifyProps) { if (groupId && props.userId) { await op.setGroup(groupId); } + markIdentitySent(key); - return { - track: (opts: EventProps & TrackProperties) => { - const { name, ...rest } = opts; - return op.track(name, groupId ? { ...rest, groups: [groupId] } : rest); - }, - }; + return { track }; } +export type Analytics = Awaited>; + /** * Noop analytics for development environment */ -- 2.51.2