From fdeaf8d595bd8b6931f269e26fe2a2e31e4e0214 Mon Sep 17 00:00:00 2001 From: Adam <2363879+adamdotdevin@users.noreply.github.com> Date: Wed, 15 Jul 2026 10:45:22 -0500 Subject: [PATCH] feat(console): add analytics logger --- bun.lock | 1 + packages/console/core/package.json | 1 + packages/console/core/src/analytics.ts | 258 +++++++++++++++++++ packages/console/core/test/analytics.test.ts | 205 +++++++++++++++ 4 files changed, 465 insertions(+) create mode 100644 packages/console/core/src/analytics.ts create mode 100644 packages/console/core/test/analytics.test.ts diff --git a/bun.lock b/bun.lock index fe623969a8..ca479e5d5b 100644 --- a/bun.lock +++ b/bun.lock @@ -202,6 +202,7 @@ "@planetscale/database": "1.19.0", "aws4fetch": "1.0.20", "drizzle-orm": "catalog:", + "effect": "catalog:", "postgres": "3.4.7", "stripe": "18.0.0", "ulid": "catalog:", diff --git a/packages/console/core/package.json b/packages/console/core/package.json index 6cb447d23d..99ed051b10 100644 --- a/packages/console/core/package.json +++ b/packages/console/core/package.json @@ -13,6 +13,7 @@ "@planetscale/database": "1.19.0", "aws4fetch": "1.0.20", "drizzle-orm": "catalog:", + "effect": "catalog:", "postgres": "3.4.7", "stripe": "18.0.0", "ulid": "catalog:", diff --git a/packages/console/core/src/analytics.ts b/packages/console/core/src/analytics.ts new file mode 100644 index 0000000000..30a341a2e4 --- /dev/null +++ b/packages/console/core/src/analytics.ts @@ -0,0 +1,258 @@ +import { Context, Effect, Layer, Logger, References, Schema } from "effect" + +export namespace AnalyticsLogger { + export const EventType = Schema.Literals(["completions", "llm.error"]) + export type EventType = typeof EventType.Type + + export const BillingSource = Schema.Literals(["anonymous", "free", "byok", "subscription", "lite", "balance"]) + export type BillingSource = typeof BillingSource.Type + + export const Request = Schema.Struct({ + id: Schema.String, + sessionID: Schema.optional(Schema.String), + projectID: Schema.optional(Schema.String), + stream: Schema.Boolean, + size: Schema.optional(Schema.Int), + responseSize: Schema.optional(Schema.Int), + status: Schema.optional(Schema.Int), + }).annotate({ identifier: "AnalyticsLogger.Request" }) + export interface Request extends Schema.Schema.Type {} + + export const Model = Schema.Struct({ + id: Schema.String, + tier: Schema.Literals(["zen", "go"]), + variant: Schema.optional(Schema.String), + }).annotate({ identifier: "AnalyticsLogger.Model" }) + export interface Model extends Schema.Schema.Type {} + + export const ProviderRoute = Schema.Struct({ + id: Schema.String, + model: Schema.String, + }).annotate({ identifier: "AnalyticsLogger.ProviderRoute" }) + export interface ProviderRoute extends Schema.Schema.Type {} + + export const Provider = Schema.Struct({ + id: Schema.String, + model: Schema.String, + shallow: Schema.optional(ProviderRoute), + budgetUsage: Schema.optional(Schema.Finite), + budgetPriority: Schema.optional(Schema.Finite), + }).annotate({ identifier: "AnalyticsLogger.Provider" }) + export interface Provider extends Schema.Schema.Type {} + + export const Account = Schema.Struct({ + source: BillingSource, + workspaceID: Schema.optional(Schema.String), + userID: Schema.optional(Schema.String), + apiKeyID: Schema.optional(Schema.String), + subscription: Schema.optional(Schema.String), + }).annotate({ identifier: "AnalyticsLogger.Account" }) + export interface Account extends Schema.Schema.Type {} + + export const Geo = Schema.Struct({ + continent: Schema.optional(Schema.String), + country: Schema.optional(Schema.String), + city: Schema.optional(Schema.String), + region: Schema.optional(Schema.String), + latitude: Schema.optional(Schema.Finite), + longitude: Schema.optional(Schema.Finite), + timezone: Schema.optional(Schema.String), + }).annotate({ identifier: "AnalyticsLogger.Geo" }) + export interface Geo extends Schema.Schema.Type {} + + export const Client = Schema.Struct({ + name: Schema.optional(Schema.String), + userAgent: Schema.optional(Schema.String), + ip: Schema.optional(Schema.String), + ipPrefix: Schema.optional(Schema.String), + geo: Schema.optional(Geo), + }).annotate({ identifier: "AnalyticsLogger.Client" }) + export interface Client extends Schema.Schema.Type {} + + export const Usage = Schema.Struct({ + input: Schema.Int, + output: Schema.Int, + reasoning: Schema.Int, + cacheRead: Schema.Int, + cacheWrite5m: Schema.Int, + cacheWrite1h: Schema.Int, + }).annotate({ identifier: "AnalyticsLogger.Usage" }) + export interface Usage extends Schema.Schema.Type {} + + // Model catalog prices are USD per one million tokens. + export const Price = Schema.Struct({ + input: Schema.Finite, + output: Schema.Finite, + cacheRead: Schema.optional(Schema.Finite), + cacheWrite5m: Schema.optional(Schema.Finite), + cacheWrite1h: Schema.optional(Schema.Finite), + }).annotate({ identifier: "AnalyticsLogger.Price" }) + export interface Price extends Schema.Schema.Type {} + + export const Cost = Schema.Struct({ + inputMicrocents: Schema.Int, + outputMicrocents: Schema.Int, + cacheReadMicrocents: Schema.optional(Schema.Int), + cacheWrite5mMicrocents: Schema.optional(Schema.Int), + cacheWrite1hMicrocents: Schema.optional(Schema.Int), + totalMicrocents: Schema.Int, + }).annotate({ identifier: "AnalyticsLogger.Cost" }) + export interface Cost extends Schema.Schema.Type {} + + export const Latency = Schema.Struct({ + totalMs: Schema.optional(Schema.Finite), + firstByteMs: Schema.optional(Schema.Finite), + firstByteAt: Schema.optional(Schema.Int), + lastByteAt: Schema.optional(Schema.Int), + }).annotate({ identifier: "AnalyticsLogger.Latency" }) + export interface Latency extends Schema.Schema.Type {} + + export const Error = Schema.Struct({ + code: Schema.optional(Schema.Union([Schema.String, Schema.Int])), + llmMessage: Schema.optional(Schema.String), + response: Schema.optional(Schema.String), + type: Schema.optional(Schema.String), + message: Schema.optional(Schema.String), + cause: Schema.optional(Schema.String), + cause2: Schema.optional(Schema.String), + }).annotate({ identifier: "AnalyticsLogger.Error" }) + export interface Error extends Schema.Schema.Type {} + + export const Event = Schema.Struct({ + version: Schema.Literal(1), + type: EventType, + timestamp: Schema.String, + dataset: Schema.String, + request: Request, + model: Model, + provider: Schema.optional(Provider), + account: Schema.optional(Account), + client: Schema.optional(Client), + usage: Schema.optional(Usage), + price: Schema.optional(Price), + cost: Schema.optional(Cost), + latency: Schema.optional(Latency), + error: Schema.optional(Error), + }).annotate({ identifier: "AnalyticsLogger.Event" }) + export interface Event extends Schema.Schema.Type {} + + type Field = string | number | boolean + export type Fields = Readonly> + export type Writer = Logger.Logger + + class Message extends Schema.TaggedClass()("AnalyticsLogger.Message", { + event: Event, + }) {} + + const isMessage = Schema.is(Message) + + export function fields(event: Event): Fields { + const cacheWriteMicrocents = + event.cost?.cacheWrite5mMicrocents === undefined && event.cost?.cacheWrite1hMicrocents === undefined + ? undefined + : (event.cost.cacheWrite5mMicrocents ?? 0) + (event.cost.cacheWrite1hMicrocents ?? 0) + const values: Record = { + _datalake_key: "inference.event", + event_version: event.version, + event_timestamp: event.timestamp, + event_date: event.timestamp.slice(0, 10), + event_type: event.type, + dataset: event.dataset, + is_stream: event.request.stream, + session: event.request.sessionID, + project: event.request.projectID, + request: event.request.id, + request_length: event.request.size, + response_length: event.request.responseSize, + status: event.request.status, + client: event.client?.name, + user_agent: event.client?.userAgent, + ip: event.client?.ip, + "ip.prefix": event.client?.ipPrefix, + "cf.continent": event.client?.geo?.continent, + "cf.country": event.client?.geo?.country, + "cf.city": event.client?.geo?.city, + "cf.region": event.client?.geo?.region, + "cf.latitude": event.client?.geo?.latitude, + "cf.longitude": event.client?.geo?.longitude, + "cf.timezone": event.client?.geo?.timezone, + model: event.model.id, + "model.tier": event.model.tier, + "model.variant": event.model.variant, + source: event.account?.source, + workspace: event.account?.workspaceID, + user_id: event.account?.userID, + api_key: event.account?.apiKeyID, + subscription: event.account?.subscription, + provider: event.provider?.id, + "provider.model": event.provider?.model, + shallowProvider: event.provider?.shallow?.id, + "shallowProvider.model": event.provider?.shallow?.model, + "provider.budget_usage": event.provider?.budgetUsage, + "provider.budget_priority": event.provider?.budgetPriority, + duration: event.latency?.totalMs, + time_to_first_byte: event.latency?.firstByteMs, + "timestamp.first_byte": event.latency?.firstByteAt, + "timestamp.last_byte": event.latency?.lastByteAt, + "tokens.input": event.usage?.input, + "tokens.output": event.usage?.output, + "tokens.reasoning": event.usage?.reasoning, + "tokens.cache_read": event.usage?.cacheRead, + "tokens.cache_write_5m": event.usage?.cacheWrite5m, + "tokens.cache_write_1h": event.usage?.cacheWrite1h, + "price.unit": event.price ? "usd_per_million_tokens" : undefined, + "price.input": event.price?.input, + "price.output": event.price?.output, + "price.cache_read": event.price?.cacheRead, + "price.cache_write_5m": event.price?.cacheWrite5m, + "price.cache_write_1h": event.price?.cacheWrite1h, + "cost.input.microcents": event.cost?.inputMicrocents, + "cost.output.microcents": event.cost?.outputMicrocents, + "cost.cache_read.microcents": event.cost?.cacheReadMicrocents, + "cost.cache_write.microcents": cacheWriteMicrocents, + "cost.cache_write_5m.microcents": event.cost?.cacheWrite5mMicrocents, + "cost.cache_write_1h.microcents": event.cost?.cacheWrite1hMicrocents, + "cost.total.microcents": event.cost?.totalMicrocents, + "llm.error.code": event.error?.code, + "llm.error.message": event.error?.llmMessage, + "error.response": event.error?.response, + "error.type": event.error?.type, + "error.message": event.error?.message, + "error.cause": event.error?.cause, + "error.cause2": event.error?.cause2, + } + return Object.fromEntries( + Object.entries(values).filter((entry): entry is [string, Field] => entry[1] !== undefined), + ) + } + + export function writer(write: (fields: Fields) => void): Writer { + return Logger.make((options) => { + const message = + Array.isArray(options.message) && options.message.length === 1 ? options.message[0] : options.message + if (!isMessage(message)) return + write(fields(message.event)) + }) + } + + export interface Interface { + readonly write: (event: Event) => Effect.Effect + } + + export class Service extends Context.Service()("@opencode/console/AnalyticsLogger") {} + + export function make(writers: ReadonlyArray): Interface { + const isolated = Logger.layer(writers, { mergeWithExisting: false }) + return Service.of({ + write: (event) => + Effect.logInfo(new Message({ event })).pipe( + Effect.provide(isolated), + Effect.provideService(References.MinimumLogLevel, "All"), + ), + }) + } + + export const layer = (writers: ReadonlyArray) => Layer.succeed(Service, make(writers)) + + export const write = (event: Event) => Effect.flatMap(Service, (service) => service.write(event)) +} diff --git a/packages/console/core/test/analytics.test.ts b/packages/console/core/test/analytics.test.ts new file mode 100644 index 0000000000..75d3ee2d7b --- /dev/null +++ b/packages/console/core/test/analytics.test.ts @@ -0,0 +1,205 @@ +import { describe, expect, test } from "bun:test" +import { Effect, Logger, References } from "effect" +import { AnalyticsLogger } from "../src/analytics" + +const event: AnalyticsLogger.Event = { + version: 1, + type: "completions", + timestamp: "2026-07-15T15:30:00.000Z", + dataset: "zen", + request: { + id: "request_01", + sessionID: "session_01", + projectID: "project_01", + stream: true, + size: 1_024, + responseSize: 4_096, + status: 200, + }, + model: { + id: "claude-opus-4-1", + tier: "go", + variant: "high", + }, + provider: { + id: "anthropic", + model: "claude-opus-4-1-20250805", + shallow: { + id: "gateway", + model: "claude-opus-4-1", + }, + budgetUsage: 0.42, + budgetPriority: 2, + }, + account: { + source: "subscription", + workspaceID: "workspace_01", + userID: "user_01", + apiKeyID: "key_01", + subscription: "20", + }, + client: { + name: "opencode", + userAgent: "opencode/1.18.1", + ip: "203.0.113.12", + ipPrefix: "203.0.113.12/32", + geo: { + continent: "NA", + country: "US", + city: "Chicago", + region: "Illinois", + latitude: 41.8781, + longitude: -87.6298, + timezone: "America/Chicago", + }, + }, + usage: { + input: 1_000, + output: 500, + reasoning: 100, + cacheRead: 200, + cacheWrite5m: 50, + cacheWrite1h: 25, + }, + price: { + input: 15, + output: 75, + cacheRead: 1.5, + cacheWrite5m: 18.75, + cacheWrite1h: 30, + }, + cost: { + inputMicrocents: 1_500_000, + outputMicrocents: 3_750_000, + cacheReadMicrocents: 30_000, + cacheWrite5mMicrocents: 93_750, + cacheWrite1hMicrocents: 75_000, + totalMicrocents: 5_448_750, + }, + latency: { + totalMs: 2_100.5, + firstByteMs: 320.25, + firstByteAt: 1_752_592_200_320, + lastByteAt: 1_752_592_202_100, + }, +} + +describe("AnalyticsLogger", () => { + test("preserves the inference fields needed by analytics", () => { + expect(AnalyticsLogger.fields(event)).toMatchObject({ + _datalake_key: "inference.event", + event_version: 1, + event_timestamp: "2026-07-15T15:30:00.000Z", + event_date: "2026-07-15", + event_type: "completions", + dataset: "zen", + is_stream: true, + session: "session_01", + project: "project_01", + request: "request_01", + client: "opencode", + user_agent: "opencode/1.18.1", + model: "claude-opus-4-1", + "model.tier": "go", + "model.variant": "high", + source: "subscription", + workspace: "workspace_01", + user_id: "user_01", + api_key: "key_01", + subscription: "20", + provider: "anthropic", + "provider.model": "claude-opus-4-1-20250805", + shallowProvider: "gateway", + "shallowProvider.model": "claude-opus-4-1", + duration: 2_100.5, + time_to_first_byte: 320.25, + "tokens.input": 1_000, + "tokens.output": 500, + "tokens.reasoning": 100, + "tokens.cache_read": 200, + "tokens.cache_write_5m": 50, + "tokens.cache_write_1h": 25, + "price.unit": "usd_per_million_tokens", + "price.input": 15, + "price.output": 75, + "price.cache_read": 1.5, + "price.cache_write_5m": 18.75, + "price.cache_write_1h": 30, + "cost.input.microcents": 1_500_000, + "cost.output.microcents": 3_750_000, + "cost.cache_read.microcents": 30_000, + "cost.cache_write.microcents": 168_750, + "cost.cache_write_5m.microcents": 93_750, + "cost.cache_write_1h.microcents": 75_000, + "cost.total.microcents": 5_448_750, + }) + }) + + test("preserves provider and application errors", () => { + const failure: AnalyticsLogger.Event = { + version: 1, + type: "llm.error", + timestamp: "2026-07-15T15:30:00.000Z", + dataset: "zen", + request: event.request, + model: event.model, + provider: event.provider, + error: { + code: "rate_limit_error", + llmMessage: "Too many requests", + response: '{"error":"rate limited"}', + type: "ProviderError", + message: "Provider rejected the request", + cause: "429", + cause2: '{"retry_after":5}', + }, + } + + expect(AnalyticsLogger.fields(failure)).toMatchObject({ + event_type: "llm.error", + "llm.error.code": "rate_limit_error", + "llm.error.message": "Too many requests", + "error.response": '{"error":"rate limited"}', + "error.type": "ProviderError", + "error.message": "Provider rejected the request", + "error.cause": "429", + "error.cause2": '{"retry_after":5}', + }) + }) + + test("writes to every analytical writer without reaching the Dash0 logger", async () => { + const first: AnalyticsLogger.Fields[] = [] + const second: AnalyticsLogger.Fields[] = [] + const dash0: unknown[] = [] + const dash0Logger = Logger.make((options) => dash0.push(options.message)) + + await Effect.runPromise( + Effect.gen(function* () { + yield* Effect.logInfo("ordinary application log") + yield* AnalyticsLogger.write(event) + }).pipe( + Effect.provide( + AnalyticsLogger.layer([ + AnalyticsLogger.writer((fields) => first.push(fields)), + AnalyticsLogger.writer((fields) => second.push(fields)), + ]), + ), + Effect.provide(Logger.layer([dash0Logger], { mergeWithExisting: false })), + Effect.provideService(References.MinimumLogLevel, "All"), + ), + ) + + expect(first).toEqual([AnalyticsLogger.fields(event)]) + expect(second).toEqual([AnalyticsLogger.fields(event)]) + expect(dash0).toEqual([["ordinary application log"]]) + }) + + test("cannot be disabled by the application log level", async () => { + const records: AnalyticsLogger.Fields[] = [] + const service = AnalyticsLogger.make([AnalyticsLogger.writer((fields) => records.push(fields))]) + + await Effect.runPromise(service.write(event).pipe(Effect.provideService(References.MinimumLogLevel, "None"))) + + expect(records).toEqual([AnalyticsLogger.fields(event)]) + }) +})