diff --git a/.gitignore b/.gitignore index f03cea3d1f..cccac291e3 100644 --- a/.gitignore +++ b/.gitignore @@ -75,3 +75,8 @@ apps/ios/build/ # Vitest Browser Mode failure screenshots and traces .vitest/ + +# bun run seed:demo writes its anchor here for the screenshot capture +scripts/seed-demo/.last-seed.json +# isolated screenshot stack env, generated by bun run seed:demo:env +.env.screenshots diff --git a/apps/ingest/package.json b/apps/ingest/package.json index aaa2538474..3671a7935a 100644 --- a/apps/ingest/package.json +++ b/apps/ingest/package.json @@ -2,7 +2,7 @@ "name": "@maple/ingest", "private": true, "scripts": { - "dev": "set -a && source ../../.env.local && set +a && cargo build && AUTUMN_SECRET_KEY= INGEST_PORT=${PORT:-3473} exec ./target/debug/maple-ingest", + "dev": "set -a && source ../../${MAPLE_DEV_ENV_FILE:-.env.local} && set +a && cargo build && AUTUMN_SECRET_KEY= INGEST_PORT=${PORT:-3473} exec ./target/debug/maple-ingest", "start": "cargo run --release", "build": "cargo build --release", "test": "cargo test", diff --git a/apps/landing/package.json b/apps/landing/package.json index 821dbdee56..41ac3e4111 100644 --- a/apps/landing/package.json +++ b/apps/landing/package.json @@ -4,6 +4,7 @@ "type": "module", "scripts": { "generate:og": "node scripts/generate-og.mjs", + "screenshots": "bun scripts/screenshots/capture.ts", "dev": "bun run --silent sync:cli && exec env -u CLAUDECODE astro dev --port ${PORT:-3391} --host ${HOST:-127.0.0.1}", "sync:i18n": "paraglide-js compile --project ./project.inlang --outdir ./src/paraglide --strategy url baseLocale && astro sync", "sync:cli": "mkdir -p public/cli && cp ../../scripts/install.sh public/cli/install && cp ../../scripts/uninstall.sh public/cli/uninstall", @@ -45,8 +46,10 @@ "tw-animate-css": "catalog:tailwind" }, "devDependencies": { + "@effect/platform-node": "catalog:effect", "@inlang/paraglide-js": "^2.25.4", "@types/d3-scale": "^4.0.9", + "playwright": "catalog:", "typescript": "catalog:tooling", "vitest": "catalog:" } diff --git a/apps/landing/public/screenshots/surface-errors.webp b/apps/landing/public/screenshots/surface-errors.webp index 1ffee84d11..5e97d2b9a6 100644 Binary files a/apps/landing/public/screenshots/surface-errors.webp and b/apps/landing/public/screenshots/surface-errors.webp differ diff --git a/apps/landing/public/screenshots/surface-service-detail.webp b/apps/landing/public/screenshots/surface-service-detail.webp new file mode 100644 index 0000000000..8cb71d28df Binary files /dev/null and b/apps/landing/public/screenshots/surface-service-detail.webp differ diff --git a/apps/landing/public/screenshots/surface-service-map.webp b/apps/landing/public/screenshots/surface-service-map.webp new file mode 100644 index 0000000000..d101af268e Binary files /dev/null and b/apps/landing/public/screenshots/surface-service-map.webp differ diff --git a/apps/landing/public/screenshots/surface-traces-peek.webp b/apps/landing/public/screenshots/surface-traces-peek.webp new file mode 100644 index 0000000000..1f2a90241d Binary files /dev/null and b/apps/landing/public/screenshots/surface-traces-peek.webp differ diff --git a/apps/landing/scripts/screenshots/capture.ts b/apps/landing/scripts/screenshots/capture.ts new file mode 100644 index 0000000000..a97117d73f --- /dev/null +++ b/apps/landing/scripts/screenshots/capture.ts @@ -0,0 +1,244 @@ +#!/usr/bin/env bun +/** + * `bun run screenshots`: capture the landing site's app screenshots from the + * local stack. Seed it first with `bun run seed:demo` at the repo root; the + * time ranges come from that run's anchor. + * + * Signs in with a one-shot Clerk ticket (dev instances only, like + * `bun run dev:signin`), renders at 2x, and writes webp via `cwebp`. + */ +import { NodeRuntime, NodeServices } from "@effect/platform-node" +import { Console, DateTime, Effect, FileSystem, Layer, Option, Schema } from "effect" +import { Command, Flag } from "effect/cli" +import { FetchHttpClient, HttpClient, HttpClientRequest } from "effect/http" +import { ChildProcess, ChildProcessSpawner } from "effect/process" +import { chromium, type BrowserContext, type Page } from "playwright" +import { RELABEL, SHOTS, TIMEZONE, VIEWPORT, type Shot } from "./shots" + +const ROOT = new URL("../../../../", import.meta.url).pathname +const SEED_STATE = `${ROOT}scripts/seed-demo/.last-seed.json` +const OUT_DIR = new URL("../../public/screenshots/", import.meta.url).pathname +const DEV_EMAIL = "david+clerk_test@gmail.com" +const MINUTE = 60_000 + +class CaptureError extends Schema.TaggedError()("@maple/landing/CaptureError", { + message: Schema.String, + shot: Schema.optional(Schema.String), +}) {} + +const SeedState = Schema.Struct({ anchor: Schema.DateTimeUtcFromString, orgId: Schema.NullOr(Schema.String) }) + +const pw = (message: string, run: () => Promise, shot?: string) => + Effect.tryPromise({ + try: run, + catch: (cause) => + new CaptureError({ + message: `${message}: ${String(cause)}`, + shot, + }), + }) + +// ── sign-in ───────────────────────────────────────────────────────────────── + +const ClerkUsers = Schema.Array(Schema.Struct({ id: Schema.String })) +const ClerkTicket = Schema.Struct({ token: Schema.String }) + +const mintTicket = Effect.fn("capture.mintTicket")( + function* (email: string) { + const secret = process.env.CLERK_SECRET_KEY ?? "" + if (!secret.startsWith("sk_test_")) { + return yield* new CaptureError({ + message: "CLERK_SECRET_KEY must be a development key (sk_test_…)", + }) + } + const client = (yield* HttpClient.HttpClient).pipe( + HttpClient.mapRequest(HttpClientRequest.bearerToken(secret)), + HttpClient.filterStatusOk, + ) + const users = yield* client + .get(`https://api.clerk.com/v1/users?email_address=${encodeURIComponent(email)}`) + .pipe( + Effect.flatMap((response) => response.json), + Effect.flatMap(Schema.decodeUnknownEffect(ClerkUsers)), + ) + const user = users[0] + if (!user) return yield* new CaptureError({ message: `no Clerk user ${email}` }) + const ticket = yield* client + .execute( + HttpClientRequest.post("https://api.clerk.com/v1/sign_in_tokens").pipe( + HttpClientRequest.bodyText( + JSON.stringify({ user_id: user.id, expires_in_seconds: 600 }), + "application/json", + ), + ), + ) + .pipe( + Effect.flatMap((response) => response.json), + Effect.flatMap(Schema.decodeUnknownEffect(ClerkTicket)), + ) + return ticket.token + }, + Effect.mapError((cause) => new CaptureError({ message: `Clerk sign-in ticket: ${cause.message}` })), +) + +// ── page prep ─────────────────────────────────────────────────────────────── + +const FREEZE_CSS = ` +*, *::before, *::after { transition: none !important; animation: none !important; caret-color: transparent !important; } +[data-sonner-toaster], #react-scan-root, .tsqd-parent-container { display: none !important; } +` + +const relabel = (page: Page) => + pw("relabel", () => + page.evaluate((pairs) => { + const walker = document.createTreeWalker(document.body, NodeFilter.SHOW_TEXT) + for (let node = walker.nextNode(); node; node = walker.nextNode()) { + for (const [from, to] of pairs) { + if (node.nodeValue?.includes(from)) node.nodeValue = node.nodeValue.replaceAll(from, to) + } + } + }, Object.entries(RELABEL)), + ) + +/** Network quiet and no skeletons left: the page is showing data, not a loading state. */ +const settle = (page: Page, shot: string) => + pw( + "wait for data", + async () => { + await page.waitForLoadState("networkidle") + await page.waitForFunction( + () => document.querySelectorAll('[data-slot="skeleton"], .animate-pulse').length === 0, + null, + { + timeout: 30_000, + }, + ) + await page.waitForTimeout(400) + }, + shot, + ) + +const warehouseTime = (ms: number) => new Date(ms).toISOString().slice(0, 19).replace("T", " ") + +const shotUrl = (base: string, shot: Shot, anchor: number) => { + const url = new URL(shot.route, base) + for (const [key, value] of Object.entries(shot.search ?? {})) { + // TanStack Router reads non-string search values as JSON. + url.searchParams.set(key, typeof value === "string" ? value : JSON.stringify(value)) + } + // Dev-only plan-gate bypass: the demo org never went through billing onboarding. + url.searchParams.set("quota_preview", "1") + if (shot.range) { + url.searchParams.set("startTime", warehouseTime(anchor - shot.range.fromMinutes * MINUTE)) + url.searchParams.set("endTime", warehouseTime(anchor - shot.range.toMinutes * MINUTE)) + } + return url.toString() +} + +const capture = Effect.fn("capture.shot")(function* ( + context: BrowserContext, + base: string, + shot: Shot, + anchor: number, +) { + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner + const fs = yield* FileSystem.FileSystem + const page = yield* Effect.acquireRelease( + pw("open page", () => context.newPage(), shot.id), + (page) => Effect.promise(() => page.close()), + ) + yield* pw("navigate", () => page.goto(shotUrl(base, shot, anchor)), shot.id) + yield* pw("freeze", () => page.addStyleTag({ content: FREEZE_CSS }), shot.id) + yield* settle(page, shot.id) + const setup = shot.setup + if (setup) { + yield* pw("setup", () => setup(page), shot.id) + yield* settle(page, shot.id) + } + yield* relabel(page) + + const png = yield* fs.makeTempFileScoped({ suffix: ".png" }) + yield* pw("screenshot", () => page.screenshot({ path: png }), shot.id) + const out = `${OUT_DIR}${shot.id}.webp` + yield* spawner + .string(ChildProcess.make("cwebp", ["-quiet", "-q", "90", png, "-o", out])) + .pipe( + Effect.mapError( + (cause) => new CaptureError({ message: `cwebp: ${cause.message}`, shot: shot.id }), + ), + ) + yield* Console.log(` ✓ ${shot.id} ${out.replace(ROOT, "")}`) +}, Effect.scoped) + +// ── command ───────────────────────────────────────────────────────────────── + +const command = Command.make( + "capture-screenshots", + { + only: Flag.String("only").pipe(Flag.withDescription("Comma-separated shot ids"), Flag.optional), + web: Flag.String("web").pipe( + Flag.withDescription("Local web origin"), + Flag.withDefault(process.env.MAPLE_DEV_WEB_URL ?? "https://web.localhost"), + ), + email: Flag.String("email").pipe(Flag.withDescription("Dev Clerk user"), Flag.withDefault(DEV_EMAIL)), + list: Flag.Boolean("list").pipe( + Flag.withDescription("Print each shot's URL and exit"), + Flag.withDefault(false), + ), + }, + Effect.fn("capture")(function* (flags) { + const fs = yield* FileSystem.FileSystem + const state = yield* fs.readFileString(SEED_STATE).pipe( + Effect.flatMap(Schema.decodeUnknownEffect(Schema.fromJsonString(SeedState))), + Effect.mapError( + () => + new CaptureError({ + message: `no seed state at ${SEED_STATE}. Run \`bun run seed:demo\` first.`, + }), + ), + ) + const anchor = DateTime.toEpochMillis(state.anchor) + const wanted = Option.map(flags.only, (only) => new Set(only.split(",").map((id) => id.trim()))) + const shots = SHOTS.filter((shot) => + Option.match(wanted, { onNone: () => true, onSome: (ids) => ids.has(shot.id) }), + ) + if (shots.length === 0) return yield* new CaptureError({ message: "no shots match --only" }) + + if (flags.list) { + for (const shot of shots) yield* Console.log(`${shot.id}\n ${shotUrl(flags.web, shot, anchor)}`) + return + } + + yield* Console.log( + `capturing ${shots.length} shot(s) at anchor ${DateTime.formatIso(state.anchor)} (org ${state.orgId ?? "unknown"})`, + ) + const ticket = yield* mintTicket(flags.email) + const browser = yield* Effect.acquireRelease( + pw("launch Chrome", () => chromium.launch({ channel: "chrome" })), + (browser) => Effect.promise(() => browser.close()), + ) + const context = yield* pw("new context", () => + browser.newContext({ + viewport: VIEWPORT, + deviceScaleFactor: 2, + timezoneId: TIMEZONE, + colorScheme: "dark", + reducedMotion: "reduce", + ignoreHTTPSErrors: true, + }), + ) + const signIn = yield* pw("open sign-in", () => context.newPage()) + yield* pw("sign in", async () => { + await signIn.goto(`${flags.web}/sign-in?__clerk_ticket=${encodeURIComponent(ticket)}`) + await signIn.waitForURL((url) => !url.pathname.startsWith("/sign-in"), { timeout: 30_000 }) + await signIn.close() + }) + + yield* Effect.forEach(shots, (shot) => capture(context, flags.web, shot, anchor), { discard: true }) + }, Effect.scoped), +).pipe(Command.withDescription("Capture landing screenshots from the seeded local stack")) + +Command.run(command, { version: "1.0.0" }).pipe( + Effect.provide(Layer.mergeAll(FetchHttpClient.layer, NodeServices.layer)), + NodeRuntime.runMain, +) diff --git a/apps/landing/scripts/screenshots/shots.ts b/apps/landing/scripts/screenshots/shots.ts new file mode 100644 index 0000000000..5f2ed69bf1 --- /dev/null +++ b/apps/landing/scripts/screenshots/shots.ts @@ -0,0 +1,46 @@ +/** + * Every app screenshot the landing site ships, as data. Captured by + * `capture.ts` against a local stack seeded with `bun run seed:demo`, with time + * ranges pinned to the seed's anchor so a re-capture frames the same moment. + */ +import type { Page } from "playwright" + +export interface Shot { + /** Output basename: `public/screenshots/.webp`. */ + readonly id: string + readonly route: string + readonly search?: Readonly>> + /** Minutes before the seed anchor the range starts and ends. Omit for pages with a fixed window. */ + readonly range?: { readonly fromMinutes: number; readonly toMinutes: number } + /** Clicks that reach the framed state once the page has data. */ + readonly setup?: (page: Page) => Promise +} + +/** Dev-account labels swapped for the demo's before every capture. */ +export const RELABEL = { + "test's Organization": "Acme", + "Redirect test US": "Acme", + "test test": "Ada Lovelace", +} as const + +export const VIEWPORT = { width: 1440, height: 810 } as const +export const TIMEZONE = "America/New_York" + +const lastMinutes = (minutes: number) => ({ fromMinutes: minutes, toMinutes: 0 }) + +export const SHOTS: ReadonlyArray = [ + { + id: "surface-traces-peek", + route: "/traces", + // Failed checkouts: payment-svc never owns the root span, so filter on the entry point. + search: { spanNames: ["POST /api/checkout"], hasError: true }, + range: lastMinutes(90), + setup: async (page) => { + await page.locator("tr[data-index]").first().click() + await page.waitForURL(/[?&]peek=/) + }, + }, + { id: "surface-service-map", route: "/service-map", range: lastMinutes(60) }, + { id: "surface-service-detail", route: "/services/payment-svc", range: lastMinutes(6 * 60) }, + { id: "surface-errors", route: "/errors" }, +] diff --git a/apps/landing/src/lib/features.ts b/apps/landing/src/lib/features.ts index 65aeb8875d..1bab80b846 100644 --- a/apps/landing/src/lib/features.ts +++ b/apps/landing/src/lib/features.ts @@ -312,8 +312,8 @@ export const features: Feature[] = [ constant: "ConnectionTimeout · payment-svc", plate: { src: "/screenshots/surface-errors.webp", - width: 2560, - height: 1320, + width: 2880, + height: 1620, alt: "Errors grouped by type with counts, affected services and last-seen times", }, facts: [ diff --git a/apps/web/src/hooks/use-mutation-action.ts b/apps/web/src/hooks/use-mutation-action.ts index 5917f75ac6..b258af894d 100644 --- a/apps/web/src/hooks/use-mutation-action.ts +++ b/apps/web/src/hooks/use-mutation-action.ts @@ -95,15 +95,3 @@ export function useMutationAction( ): readonly [run: (value: W) => Promise>, pending: boolean] { return useAsyncAction(useMutationRunner(atom, options)) } - -/** - * `useMutationAction` for row actions: `run(rowId, payload)` and `isPending(rowId)`, so a list - * shows the spinner on the row being mutated instead of a shared flag. - */ -export function useKeyedMutationAction( - atom: Atom.Writable, W>, - options: MutationActionOptions, -): KeyedAction> { - const runner = useMutationRunner(atom, options) - return useKeyedAsyncAction((_key: K, value: W) => runner(value)) -} diff --git a/apps/web/src/lib/error-toast.ts b/apps/web/src/lib/error-toast.ts index 4da249c6e7..2ba8f9503f 100644 --- a/apps/web/src/lib/error-toast.ts +++ b/apps/web/src/lib/error-toast.ts @@ -27,12 +27,6 @@ export const errorMessage = (error: unknown, fallback: string): string => { return isUnexpectedError(presentation) ? fallback : presentation.message } -/** The user-facing message of a failed Exit, or `fallback` for a success, a non-Exit or an unexpected defect. */ -export function getExitErrorMessage(exit: unknown, fallback: string): string { - if (!Exit.isExit(exit) || Exit.isSuccess(exit)) return fallback - return errorMessage(exit, fallback) -} - interface ToastExitOptions { /** Toast title on success; omit to stay silent on success. */ readonly success?: string diff --git a/bun.lock b/bun.lock index 353049524d..d6e316f36a 100644 --- a/bun.lock +++ b/bun.lock @@ -208,8 +208,10 @@ "tw-animate-css": "catalog:tailwind", }, "devDependencies": { + "@effect/platform-node": "catalog:effect", "@inlang/paraglide-js": "^2.25.4", "@types/d3-scale": "^4.0.9", + "playwright": "catalog:", "typescript": "catalog:tooling", "vitest": "catalog:", }, diff --git a/package.json b/package.json index cbd6643db2..e8dac5e20b 100644 --- a/package.json +++ b/package.json @@ -11,6 +11,8 @@ "build": "turbo build", "dev": "bun run ./scripts/dev.ts", "dev:signin": "bun run ./scripts/dev-signin.ts", + "seed:demo": "bun run ./scripts/seed-demo.ts", + "seed:demo:env": "bun run ./scripts/seed-demo/make-env.ts", "prepare": "effect-tsgo unpatch --no-typescript --oxlint && effect-tsgo patch --no-typescript --oxlint", "docker:build": "COMPOSE_PARALLEL_LIMIT=1 docker compose build", "docker:up": "COMPOSE_PARALLEL_LIMIT=1 docker compose up --build", diff --git a/packages/infra/src/cloudflare/maple-db.ts b/packages/infra/src/cloudflare/maple-db.ts index 8d90a50b25..a71968cdb8 100644 --- a/packages/infra/src/cloudflare/maple-db.ts +++ b/packages/infra/src/cloudflare/maple-db.ts @@ -51,13 +51,15 @@ export const ManagedMapleDb = Cloudflare.Hyperdrive.Connection( }, // Read-after-write everywhere (alert state CAS, dashboard versioning). caching: { disabled: true }, + // The origin local workers actually dial. From `MAPLE_PG_URL` too: a hardcoded + // database here silently split a stack across two databases. dev: { scheme: "postgres", - host: "localhost", - port: 5499, - database: "maple", - user: "maple", - password: Redacted.make("maple"), + host: pgUrl.hostname, + port: Number(pgUrl.port || "5432"), + database: pgUrl.pathname.replace(/^\//, "") || "postgres", + user: decodeURIComponent(pgUrl.username), + password: Redacted.make(decodeURIComponent(pgUrl.password)), // Docker Postgres has no TLS; alchemy's default `prefer` stalls until timeout. sslmode: "disable", }, diff --git a/scripts/dev.ts b/scripts/dev.ts index 26ff3317bb..2126b6244b 100644 --- a/scripts/dev.ts +++ b/scripts/dev.ts @@ -3,6 +3,7 @@ * whole stack. Everything else lives in alchemy.run.ts. */ import { spawn } from "node:child_process" +import { existsSync, readFileSync } from "node:fs" import path from "node:path" import { DEV_APPS, DEV_APPS_ENV_KEY, isDevApp, type DevApp } from "../packages/infra/src/dev-urls.ts" @@ -19,6 +20,18 @@ if (unknown.length > 0) { const selected: ReadonlyArray = args.length > 0 ? DEV_APPS.filter((app) => args.includes(app)) : DEV_APPS +// `.env.screenshots` (from `bun run seed:demo:env`) runs the isolated screenshot stack. +const envFile = process.env.MAPLE_DEV_ENV_FILE ?? ".env.local" + +// mise loads `.env.local` into the shell, and a blank value in an env file counts as +// absent, so with any other env file those keys would leak through. Drop them. +const localKeys = new Set( + envFile === ".env.local" || !existsSync(".env.local") + ? [] + : readFileSync(".env.local", "utf8").match(/^[A-Z0-9_]+(?==)/gm), +) +const inherited = Object.fromEntries(Object.entries(process.env).filter(([key]) => !localKeys.has(key))) + // A spawned child does not inherit the `node_modules/.bin` PATH entry `bun run` adds. const alchemyBin = path.join(import.meta.dirname, "..", "node_modules", ".bin", "alchemy") @@ -29,14 +42,14 @@ const child = spawn( "--stage", process.env.MAPLE_DEV_STAGE ?? `dev_${process.env.USER ?? "local"}`, "--env-file", - ".env.local", + envFile, ], { stdio: "inherit", // Own process group: a signal to alchemy's CLI process alone stops nothing. detached: true, env: { - ...process.env, + ...inherited, [DEV_APPS_ENV_KEY]: selected.join(","), // Dev stacks never touch the shared account state store. ALCHEMY_LOCAL_STATE: process.env.ALCHEMY_LOCAL_STATE ?? "1", diff --git a/scripts/seed-demo.ts b/scripts/seed-demo.ts new file mode 100644 index 0000000000..688d8b3ef4 --- /dev/null +++ b/scripts/seed-demo.ts @@ -0,0 +1,363 @@ +#!/usr/bin/env bun +/** + * `bun run seed:demo`: fill the local stack with the Acme Shop world (see + * `seed-demo/scenario.ts`) for landing screenshots and UI checks. + * + * Telemetry goes through the local ingest gateway like real traffic, so it + * lands in the org `MAPLE_ORG_ID_OVERRIDE` names. The run's time anchor is + * written to `seed-demo/.last-seed.json`; screenshots pin their time ranges to it. + */ +import { NodeRuntime, NodeServices } from "@effect/platform-node" +import { Console, DateTime, Effect, FileSystem, Layer, Option, Schema } from "effect" +import { Command, Flag } from "effect/cli" +import { FetchHttpClient, HttpClient, HttpClientRequest } from "effect/http" +import { ChildProcess, ChildProcessSpawner } from "effect/process" +import * as Schedule from "effect/Schedule" +import { resourceLogs, resourceMetrics, resourceSpans } from "./seed-demo/otlp" +import { Rng } from "./seed-demo/rng" +import { generateWindow, World, type Grouped, type WindowOutput } from "./seed-demo/telemetry" + +const HOUR = 3_600_000 +const MINUTE = 60_000 +const SPANS_PER_REQUEST = 2_000 +const STATE_FILE = new URL("./seed-demo/.last-seed.json", import.meta.url).pathname + +class SeedPreflightError extends Schema.TaggedError()("@maple/seed-demo/PreflightError", { + message: Schema.String, +}) {} + +const SeedState = Schema.Struct({ + anchor: Schema.DateTimeUtcFromString, + start: Schema.DateTimeUtcFromString, + incidentAt: Schema.optionalKey(Schema.DateTimeUtcFromString), + seed: Schema.optionalKey(Schema.Number), + followedThrough: Schema.optionalKey(Schema.DateTimeUtcFromString), +}) + +class SeedIngestError extends Schema.TaggedError()("@maple/seed-demo/IngestError", { + message: Schema.String, + path: Schema.String, +}) {} + +interface Target { + readonly endpoint: string + readonly key: string +} + +const env = (name: string) => Option.fromNullishOr(process.env[name]?.trim() || undefined) + +// ── preflight ─────────────────────────────────────────────────────────────── + +const checkIngest = Effect.fn("seedDemo.checkIngest")(function* (target: Target) { + const client = yield* HttpClient.HttpClient + const response = yield* client.execute(HttpClientRequest.get(`${target.endpoint}/health`)).pipe( + Effect.mapError( + () => + new SeedPreflightError({ + message: `the ingest gateway is not answering at ${target.endpoint}. Start it with \`bun dev ingest\`.`, + }), + ), + ) + if (response.status !== 200) { + return yield* new SeedPreflightError({ + message: `the ingest gateway at ${target.endpoint}/health answered ${response.status}`, + }) + } +}) + +/** A BYO ClickHouse row makes seeded data invisible and puts that warehouse's data in the shots. */ +const checkWarehouseRouting = Effect.fn("seedDemo.checkWarehouseRouting")(function* (orgId: string) { + const pgUrl = env("MAPLE_PG_URL") + if (Option.isNone(pgUrl) || !/^org_[A-Za-z0-9]+$/.test(orgId)) { + yield* Console.warn(" ! skipped the BYO ClickHouse check (no MAPLE_PG_URL, or an unexpected org id)") + return + } + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner + const output = yield* spawner + .string( + ChildProcess.make("psql", [ + pgUrl.value, + "-tAc", + `SELECT ch_url FROM org_clickhouse_settings WHERE org_id = '${orgId}'`, + ]), + ) + .pipe(Effect.option) + if (Option.isNone(output)) { + yield* Console.warn(" ! could not run psql; skipped the BYO ClickHouse check") + return + } + const chUrl = output.value.trim() + if (chUrl !== "") { + return yield* new SeedPreflightError({ + message: + `org ${orgId} reads from its own ClickHouse (${chUrl}), so seeded data would be invisible and ` + + "screenshots would show that warehouse instead. Remove it with DELETE /api/org-clickhouse-settings/ " + + "(not raw SQL: the config is edge-cached for an hour) and re-run.", + }) + } +}) + +const checkWarehouse = Effect.fn("seedDemo.checkWarehouse")(function* () { + const host = env("TINYBIRD_HOST") + if (Option.isNone(host)) return + const client = yield* HttpClient.HttpClient + const reachable = yield* client.execute(HttpClientRequest.get(host.value)).pipe(Effect.option) + if (Option.isNone(reachable)) { + yield* Console.warn( + ` ! TINYBIRD_HOST (${host.value}) is not answering. Ingest will accept the data, but nothing lands until it is up.`, + ) + } +}) + +// ── sending ───────────────────────────────────────────────────────────────── + +const post = Effect.fn("seedDemo.post")( + function* (target: Target, path: string, body: unknown) { + const client = yield* HttpClient.HttpClient + const request = HttpClientRequest.post(`${target.endpoint}${path}`).pipe( + HttpClientRequest.bearerToken(target.key), + HttpClientRequest.bodyText(JSON.stringify(body), "application/json"), + ) + const response = yield* client + .execute(request) + .pipe(Effect.mapError((cause) => new SeedIngestError({ message: cause.message, path }))) + if (response.status >= 300) { + const text = yield* response.text.pipe(Effect.orElseSucceed(() => "")) + return yield* new SeedIngestError({ message: `${response.status}: ${text.slice(0, 300)}`, path }) + } + }, + Effect.retry({ times: 3, schedule: Schedule.exponential("250 millis") }), +) + +const chunk = (items: ReadonlyArray, size: number): ReadonlyArray> => + Array.from({ length: Math.ceil(items.length / size) }, (_, i) => items.slice(i * size, (i + 1) * size)) + +const requestsFor = (out: WindowOutput) => [ + ...out.spans.flatMap((group) => + chunk(group.items, SPANS_PER_REQUEST).map((spans) => ({ + path: "/v1/traces", + body: resourceSpans(group.resource.attributes, spans), + })), + ), + ...out.logs.map((group) => ({ + path: "/v1/logs", + body: resourceLogs(group.resource.attributes, group.items), + })), + ...out.metrics.map((group) => ({ + path: "/v1/metrics", + body: resourceMetrics(group.resource.attributes, group.items), + })), +] + +const count = (groups: ReadonlyArray>) => + groups.reduce((total, group) => total + group.items.length, 0) + +const clock = (ms: number) => new Date(ms).toTimeString().slice(0, 5) + +// ── command ───────────────────────────────────────────────────────────────── + +const seed = Command.make( + "seed-demo", + { + endpoint: Flag.String("endpoint").pipe( + Flag.withDescription("Ingest gateway base URL"), + Flag.withDefault(Option.getOrElse(env("MAPLE_INGEST_URL"), () => "http://localhost:3474")), + ), + key: Flag.String("key").pipe( + Flag.withDescription("Ingest key for the org to seed (env MAPLE_INGEST_KEY)"), + Flag.withDefault(Option.getOrElse(env("MAPLE_INGEST_KEY"), () => "maple_pk_local")), + ), + hours: Flag.Int("hours").pipe( + Flag.withDescription("How much history to backfill"), + Flag.withDefault(24), + ), + rate: Flag.Int("rate").pipe(Flag.withDescription("Peak traces per minute"), Flag.withDefault(60)), + incidentMinutes: Flag.Int("incident-minutes").pipe( + Flag.withDescription("How long before the anchor the bad payment-svc deploy lands"), + Flag.withDefault(120), + ), + anchor: Flag.String("anchor").pipe( + Flag.withDescription( + 'End of the seeded window: "now" or ISO-8601 (default: now, floored to 5 minutes)', + ), + Flag.withDefault("now"), + ), + seed: Flag.Int("seed").pipe( + Flag.withDescription("RNG seed; same seed, same data"), + Flag.withDefault(7), + ), + dryRun: Flag.Boolean("dry-run").pipe( + Flag.withDescription("Generate and count, send nothing"), + Flag.withDefault(false), + ), + follow: Flag.Boolean("follow").pipe( + Flag.withDescription( + "After the backfill, keep sending each minute live (the error tick only scans active orgs)", + ), + Flag.withDefault(false), + ), + resume: Flag.Boolean("resume").pipe( + Flag.withDescription("Skip the backfill and follow on from where the last run stopped sending"), + Flag.withDefault(false), + ), + }, + Effect.fn("seedDemo")(function* (flags) { + const fs = yield* FileSystem.FileSystem + const previous = flags.resume + ? Option.some( + yield* fs.readFileString(STATE_FILE).pipe( + Effect.flatMap(Schema.decodeUnknownEffect(Schema.fromJsonString(SeedState))), + Effect.mapError( + () => + new SeedPreflightError({ + message: `--resume: no readable state at ${STATE_FILE}`, + }), + ), + ), + ) + : Option.none() + const target: Target = { endpoint: flags.endpoint.replace(/\/$/, ""), key: flags.key } + const anchor = Option.match(previous, { + onSome: (state) => DateTime.toEpochMillis(state.anchor), + onNone: () => + flags.anchor === "now" + ? Math.floor(Date.now() / (5 * MINUTE)) * 5 * MINUTE + : Option.match(Schema.decodeUnknownOption(Schema.DateTimeUtcFromString)(flags.anchor), { + onSome: DateTime.toEpochMillis, + onNone: () => Number.NaN, + }), + }) + const hours = Option.isSome(previous) ? 0 : flags.hours + if (Number.isNaN(anchor)) { + return yield* new SeedPreflightError({ message: `--anchor: could not parse "${flags.anchor}"` }) + } + const orgId = env("MAPLE_ORG_ID_OVERRIDE") + + if (!flags.dryRun) { + yield* Console.log("preflight") + yield* checkIngest(target) + if (Option.isNone(orgId)) { + yield* Console.warn( + " ! MAPLE_ORG_ID_OVERRIDE is unset; data lands in whatever org the ingest key resolves to", + ) + } else { + yield* checkWarehouseRouting(orgId.value) + yield* Console.log(` org ${orgId.value}`) + } + yield* checkWarehouse() + } + + // A resumed run reuses the saved seed and incident, so pods and the incident line up. + const seed = Option.match(previous, { + onSome: (state) => state.seed ?? flags.seed, + onNone: () => flags.seed, + }) + const incidentBeforeAnchorMs = Option.match(previous, { + onSome: (state) => + state.incidentAt === undefined + ? flags.incidentMinutes * MINUTE + : anchor - DateTime.toEpochMillis(state.incidentAt), + onNone: () => flags.incidentMinutes * MINUTE, + }) + const world = new World({ + anchor, + // A resumed run keeps the original window, so pods, versions and the incident line up. + windowMs: Option.match(previous, { + onSome: (state) => anchor - DateTime.toEpochMillis(state.start), + onNone: () => hours * HOUR, + }), + incidentBeforeAnchorMs, + peakTracesPerMinute: flags.rate, + seed, + }) + yield* Console.log( + `\nseeding ${clock(world.start)} → ${clock(anchor)} (${hours}h), incident at ${clock(world.incidentAt)}` + + (flags.dryRun ? " [dry run]" : ""), + ) + + const totals = { traces: 0, errors: 0, spans: 0, logs: 0 } + // Seeded per window, so a follow run never replays the backfill's trace ids. + const rngFor = (from: number) => new Rng(Math.imul(seed, 0x9e3779b1) ^ Math.floor(from / MINUTE)) + const sendWindow = (from: number, to: number) => + Effect.gen(function* () { + const out = generateWindow(world, rngFor(from), from, to) + const spans = count(out.spans) + const logs = count(out.logs) + totals.traces += out.traceCount + totals.errors += out.errorTraceCount + totals.spans += spans + totals.logs += logs + if (!flags.dryRun) { + yield* Effect.forEach(requestsFor(out), ({ path, body }) => post(target, path, body), { + concurrency: 4, + discard: true, + }) + } + yield* Console.log( + ` ${clock(from)} ${String(out.traceCount).padStart(6)} traces ${String(out.errorTraceCount).padStart(4)} failed ` + + `${String(spans).padStart(7)} spans ${String(logs).padStart(6)} logs`, + ) + }) + const windows = Array.from({ length: hours }, (_, i) => anchor - hours * HOUR + i * HOUR) + yield* Effect.forEach(windows, (from) => sendWindow(from, Math.min(from + HOUR, anchor))) + + yield* Console.log( + `\n${flags.dryRun ? "would send" : "sent"} ${totals.traces} traces (${totals.errors} failed), ${totals.spans} spans, ${totals.logs} logs`, + ) + if (flags.dryRun) return + + const writeState = (followedThrough: number) => + fs.writeFileString( + STATE_FILE, + `${JSON.stringify( + { + anchor: new Date(anchor).toISOString(), + start: new Date(world.start).toISOString(), + incidentAt: new Date(world.incidentAt).toISOString(), + followedThrough: new Date(followedThrough).toISOString(), + orgId: Option.getOrNull(orgId), + seed, + }, + null, + "\t", + )}\n`, + ) + let cursor = Option.match(previous, { + onSome: (state) => DateTime.toEpochMillis(state.followedThrough ?? state.anchor), + onNone: () => anchor, + }) + yield* writeState(cursor) + yield* Console.log( + `anchor written to ${STATE_FILE}\nData arrives through the pipeline; give it a minute to appear.`, + ) + if (!flags.follow && !flags.resume) return + + yield* Console.log(`\nfollowing from ${clock(cursor)}: one window per minute, Ctrl-C to stop`) + yield* Effect.gen(function* () { + const end = Math.floor(Date.now() / MINUTE) * MINUTE + const catchUp = Array.from( + { length: Math.ceil((end - cursor) / HOUR) }, + (_, i) => cursor + i * HOUR, + ) + yield* Effect.forEach(catchUp, (from) => + Effect.gen(function* () { + const to = Math.min(from + HOUR, end) + yield* sendWindow(from, to) + cursor = to + yield* writeState(cursor) + }), + ).pipe( + // A restarting gateway is not fatal: the unsent windows go again next minute. + Effect.catchTag("@maple/seed-demo/IngestError", (error) => + Console.warn(` ! ${error.path}: ${error.message}; retrying next minute`), + ), + ) + yield* Effect.sleep(`${MINUTE - (Date.now() % MINUTE) + 2_000} millis`) + }).pipe(Effect.forever) + }), +).pipe(Command.withDescription("Seed the local stack with the Acme Shop demo world")) + +Command.run(seed, { version: "1.0.0" }).pipe( + Effect.provide(Layer.mergeAll(FetchHttpClient.layer, NodeServices.layer)), + NodeRuntime.runMain, +) diff --git a/scripts/seed-demo/make-env.ts b/scripts/seed-demo/make-env.ts new file mode 100644 index 0000000000..f0ec8dcafa --- /dev/null +++ b/scripts/seed-demo/make-env.ts @@ -0,0 +1,100 @@ +#!/usr/bin/env bun +/** + * `bun run seed:demo:env`: set up the isolated screenshot stack. + * + * Creates and migrates a `maple_screenshots` Postgres database next to the dev + * one, then writes `.env.screenshots`: `.env.local` pinned to the demo org, on + * that database, with the static ingest key store and the alerting crons on. + * The crons only ever see the demo org's rules, so the seeded incident never + * pages a real destination. Run the stack with + * `MAPLE_DEV_ENV_FILE=.env.screenshots bun dev api web ingest alerting`. + */ +import { NodeRuntime, NodeServices } from "@effect/platform-node" +import { Console, Effect, FileSystem, Schema } from "effect" +import { ChildProcess, ChildProcessSpawner } from "effect/process" + +export const DEMO_ORG_ID = "org_3AcmeShopDemo0000000000000" +const DATABASE = "maple_screenshots" +const ROOT = new URL("../../", import.meta.url).pathname + +class MakeEnvError extends Schema.TaggedError()("@maple/seed-demo/MakeEnvError", { + message: Schema.String, +}) {} + +// Explicit values only: mise loads `.env.local` into the shell, so a dropped key would still be inherited. +const OVERRIDES = new Map([ + ["MAPLE_ORG_ID_OVERRIDE", DEMO_ORG_ID], + ["MAPLE_ALERTING_ALLOW_NONPROD", "1"], + // Any maple_pk_* key resolves to the demo org. + ["INGEST_KEY_STORE_BACKEND", "static"], + // Billing off: the plan gate fails open and no test customer lands in the Autumn sandbox. + ["AUTUMN_SECRET_KEY", ""], +]) + +/** Point the URL at the screenshot database, also when it names none or ends in a slash. */ +const withDatabase = (url: string) => + Effect.try({ + try: () => { + const parsed = new URL(url) + parsed.pathname = `/${DATABASE}` + return parsed.toString() + }, + catch: () => new MakeEnvError({ message: "MAPLE_PG_URL in .env.local is not a valid URL" }), + }) + +const program = Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner + const source = yield* fs.readFileString(`${ROOT}.env.local`) + + const pgLine = source.split("\n").find((line) => line.startsWith("MAPLE_PG_URL=")) + if (pgLine === undefined) return yield* new MakeEnvError({ message: ".env.local has no MAPLE_PG_URL" }) + const devUrl = pgLine.slice("MAPLE_PG_URL=".length).trim() + const demoUrl = yield* withDatabase(devUrl) + + const exists = yield* spawner.string( + ChildProcess.make("psql", [ + devUrl, + "-tAc", + `SELECT 1 FROM pg_database WHERE datname = '${DATABASE}'`, + ]), + ) + if (exists.trim() !== "1") { + yield* spawner.string(ChildProcess.make("psql", [devUrl, "-c", `CREATE DATABASE ${DATABASE}`])) + yield* Console.log(`created database ${DATABASE}`) + } + yield* spawner.string( + ChildProcess.make("bun", ["run", "db:migrate"], { + cwd: `${ROOT}packages/db`, + env: { DATABASE_URL: demoUrl }, + extendEnv: true, + }), + { includeStderr: true }, + ) + yield* Console.log(`migrated ${DATABASE}`) + + const seen = new Set() + const lines = source.split("\n").flatMap((line) => { + const key = /^([A-Z0-9_]+)=/.exec(line)?.[1] + if (key === "MAPLE_PG_URL") return [`MAPLE_PG_URL=${demoUrl}`] + if (key === undefined || !OVERRIDES.has(key)) return [line] + seen.add(key) + return [`${key}=${OVERRIDES.get(key)}`] + }) + const added = [...OVERRIDES].flatMap(([key, value]) => (seen.has(key) ? [] : [`${key}=${value}`])) + // Copies .env.local's secrets, so owner-only; chmod covers a file an older run created. + const envPath = `${ROOT}.env.screenshots` + yield* fs.writeFileString( + envPath, + [ + `# Generated by bun run seed:demo:env from .env.local. Do not edit; re-run instead.`, + ...lines, + ...added, + ].join("\n"), + { mode: 0o600 }, + ) + yield* fs.chmod(envPath, 0o600) + yield* Console.log(`wrote .env.screenshots (org ${DEMO_ORG_ID}, database ${DATABASE})`) +}) + +program.pipe(Effect.provide(NodeServices.layer), NodeRuntime.runMain) diff --git a/scripts/seed-demo/otlp.ts b/scripts/seed-demo/otlp.ts new file mode 100644 index 0000000000..a681211ff9 --- /dev/null +++ b/scripts/seed-demo/otlp.ts @@ -0,0 +1,146 @@ +/** + * OTLP/JSON shapes, as the ingest gateway decodes them (opentelemetry-proto + * `with-serde`): camelCase keys, hex ids, `*UnixNano` as decimal strings. + */ +export type AnyValue = + | { stringValue: string } + | { intValue: string } + | { doubleValue: number } + | { boolValue: boolean } +export interface KeyValue { + readonly key: string + readonly value: AnyValue +} + +export const attr = (key: string, value: string | number | boolean): KeyValue => { + if (typeof value === "boolean") return { key, value: { boolValue: value } } + if (typeof value === "string") return { key, value: { stringValue: value } } + return Number.isInteger(value) + ? { key, value: { intValue: String(value) } } + : { key, value: { doubleValue: value } } +} + +export const nano = (ms: number): string => (BigInt(Math.round(ms * 1000)) * 1000n).toString() + +export const SpanKind = { internal: 1, server: 2, client: 3, producer: 4, consumer: 5 } as const +export const StatusCode = { unset: 0, error: 2 } as const +export const Severity = { + DEBUG: 5, + INFO: 9, + WARN: 13, + ERROR: 17, +} as const +export type SeverityText = keyof typeof Severity + +export interface SpanEvent { + readonly timeUnixNano: string + readonly name: string + readonly attributes: ReadonlyArray +} + +export interface Span { + readonly traceId: string + readonly spanId: string + readonly parentSpanId?: string | undefined + readonly name: string + readonly kind: number + readonly startTimeUnixNano: string + readonly endTimeUnixNano: string + readonly attributes: ReadonlyArray + /** `undefined` fields drop out of the JSON body. */ + readonly events?: ReadonlyArray | undefined + readonly status: { readonly code: number; readonly message?: string } +} + +export interface LogRecord { + readonly timeUnixNano: string + readonly observedTimeUnixNano: string + readonly severityNumber: number + readonly severityText: SeverityText + readonly body: { readonly stringValue: string } + readonly attributes: ReadonlyArray + readonly traceId: string | undefined + readonly spanId: string | undefined +} + +export interface NumberPoint { + readonly startTimeUnixNano?: string + readonly timeUnixNano: string + readonly attributes: ReadonlyArray + readonly asDouble?: number + readonly asInt?: number +} + +export interface HistogramPoint { + readonly startTimeUnixNano: string + readonly timeUnixNano: string + readonly attributes: ReadonlyArray + readonly count: string + readonly sum: number + readonly min: number + readonly max: number + readonly bucketCounts: ReadonlyArray + readonly explicitBounds: ReadonlyArray +} + +const DELTA = 1 + +export type Metric = + | { + readonly name: string + readonly unit: string + readonly gauge: { readonly dataPoints: ReadonlyArray } + } + | { + readonly name: string + readonly unit: string + readonly sum: { + readonly aggregationTemporality: typeof DELTA + readonly isMonotonic: boolean + readonly dataPoints: ReadonlyArray + } + } + | { + readonly name: string + readonly unit: string + readonly histogram: { + readonly aggregationTemporality: typeof DELTA + readonly dataPoints: ReadonlyArray + } + } + +export const gauge = (name: string, unit: string, dataPoints: ReadonlyArray): Metric => ({ + name, + unit, + gauge: { dataPoints }, +}) + +export const deltaSum = (name: string, unit: string, dataPoints: ReadonlyArray): Metric => ({ + name, + unit, + sum: { aggregationTemporality: DELTA, isMonotonic: true, dataPoints }, +}) + +export const deltaHistogram = ( + name: string, + unit: string, + dataPoints: ReadonlyArray, +): Metric => ({ + name, + unit, + histogram: { aggregationTemporality: DELTA, dataPoints }, +}) + +const SCOPE = { name: "maple-seed-demo", version: "1.0.0" } + +export const resourceSpans = (resource: ReadonlyArray, spans: ReadonlyArray) => ({ + resourceSpans: [{ resource: { attributes: resource }, scopeSpans: [{ scope: SCOPE, spans }] }], +}) + +export const resourceLogs = (resource: ReadonlyArray, logRecords: ReadonlyArray) => ({ + resourceLogs: [{ resource: { attributes: resource }, scopeLogs: [{ scope: SCOPE, logRecords }] }], +}) + +export const resourceMetrics = (resource: ReadonlyArray, metrics: ReadonlyArray) => ({ + resourceMetrics: [{ resource: { attributes: resource }, scopeMetrics: [{ scope: SCOPE, metrics }] }], +}) diff --git a/scripts/seed-demo/rng.ts b/scripts/seed-demo/rng.ts new file mode 100644 index 0000000000..efcd3e48bf --- /dev/null +++ b/scripts/seed-demo/rng.ts @@ -0,0 +1,53 @@ +/** Seeded mulberry32. Ids come from it too, so a re-seed reproduces the same traces. */ +export class Rng { + private state: number + + constructor(seed: number) { + this.state = seed >>> 0 + } + + next(): number { + this.state = (this.state + 0x6d2b79f5) | 0 + let t = Math.imul(this.state ^ (this.state >>> 15), 1 | this.state) + t = (t + Math.imul(t ^ (t >>> 7), 61 | t)) ^ t + return ((t ^ (t >>> 14)) >>> 0) / 4294967296 + } + + int(min: number, max: number): number { + return min + Math.floor(this.next() * (max - min + 1)) + } + + chance(p: number): boolean { + return this.next() < p + } + + pick(items: readonly [T, ...T[]]): T { + return items[Math.floor(this.next() * items.length)] ?? items[0] + } + + weighted(items: readonly [readonly [T, number], ...(readonly [T, number])[]]): T { + const total = items.reduce((sum, [, weight]) => sum + weight, 0) + let roll = this.next() * total + for (const [item, weight] of items) { + roll -= weight + if (roll <= 0) return item + } + return items[0][0] + } + + gaussian(): number { + const u = Math.max(this.next(), Number.EPSILON) + return Math.sqrt(-2 * Math.log(u)) * Math.cos(2 * Math.PI * this.next()) + } + + /** Log-normal latency: `median` ms, `spread` ≈ 0.3 tight, 0.8 long-tailed. */ + latency(median: number, spread: number): number { + return Math.max(0.2, median * Math.exp(spread * this.gaussian())) + } + + hex(bytes: number): string { + let out = "" + for (let i = 0; i < bytes; i++) out += this.int(0, 255).toString(16).padStart(2, "0") + return out + } +} diff --git a/scripts/seed-demo/scenario.ts b/scripts/seed-demo/scenario.ts new file mode 100644 index 0000000000..5c8394a010 --- /dev/null +++ b/scripts/seed-demo/scenario.ts @@ -0,0 +1,515 @@ +/** + * Acme Shop: the one world every landing screenshot is taken in. + * + * Nine services on Kubernetes, 24h of traffic on a daily curve, a handful of + * routine deploys, and one bad one: payment-svc 3.5.0 lands `incidentOffsetMs` + * before the anchor, exhausts its Postgres pool, and starts timing out charges. + */ +import type { Rng } from "./rng" + +export type Language = "nodejs" | "go" | "python" + +export type ServiceName = + | "storefront" + | "checkout-api" + | "cart-svc" + | "catalog-api" + | "inventory-svc" + | "payment-svc" + | "auth-svc" + | "order-worker" + | "notification-svc" + +export interface ServiceDef { + readonly name: ServiceName + readonly language: Language + readonly baseVersion: string + /** Pods per rollout; each span lands on one of them. */ + readonly replicas: number +} + +export const SERVICES: ReadonlyArray = [ + { name: "storefront", language: "nodejs", baseVersion: "5.8.2", replicas: 4 }, + { name: "checkout-api", language: "nodejs", baseVersion: "4.2.0", replicas: 3 }, + { name: "cart-svc", language: "go", baseVersion: "1.22.1", replicas: 2 }, + { name: "catalog-api", language: "python", baseVersion: "2.14.0", replicas: 3 }, + { name: "inventory-svc", language: "go", baseVersion: "3.1.4", replicas: 2 }, + { name: "payment-svc", language: "nodejs", baseVersion: "3.4.1", replicas: 2 }, + { name: "auth-svc", language: "go", baseVersion: "2.0.3", replicas: 2 }, + { name: "order-worker", language: "nodejs", baseVersion: "1.9.0", replicas: 2 }, + { name: "notification-svc", language: "python", baseVersion: "0.12.5", replicas: 1 }, +] + +export interface DeployDef { + readonly service: ServiceName + readonly version: string + /** Milliseconds before the anchor. */ + readonly beforeAnchorMs: number +} + +const HOUR = 3_600_000 +const MINUTE = 60_000 + +export const DEPLOYS: ReadonlyArray = [ + { service: "catalog-api", version: "2.15.0", beforeAnchorMs: 19 * HOUR }, + { service: "storefront", version: "5.9.0", beforeAnchorMs: 10 * HOUR }, + { service: "cart-svc", version: "1.23.0", beforeAnchorMs: 5 * HOUR + 30 * MINUTE }, + { service: "payment-svc", version: "3.5.0", beforeAnchorMs: 2 * HOUR }, +] + +export const INCIDENT = { + service: "payment-svc", + version: "3.5.0", + poolSize: 20, +} as const satisfies { readonly service: ServiceName; readonly version: string; readonly poolSize: number } + +/** What an op can see about the moment its trace runs in. */ +export interface Ctx { + readonly at: number + /** True once the bad payment-svc deploy is live. */ + readonly incident: boolean + readonly rng: Rng +} + +export interface Failure { + readonly type: string + readonly message: string + readonly stacktrace: string + /** Log line lead-in, e.g. "charge failed"; defaults to the exception type. */ + readonly logPrefix?: string + /** The span runs this long when it fails (a timeout ends at the timeout). */ + readonly durationMs?: number +} + +type Latency = (ctx: Ctx) => readonly [median: number, spread: number] +type FailRule = (ctx: Ctx) => Failure | null + +export type Op = + | { + readonly kind: "server" + readonly service: ServiceName + readonly method: string + readonly route: string + readonly self: Latency + readonly children: ReadonlyArray + readonly fail?: FailRule | undefined + } + | { + readonly kind: "db" + readonly system: "postgresql" | "redis" + readonly namespace: string + readonly operation: string + readonly collection?: string + readonly statement: string + readonly latency: Latency + readonly repeat?: (ctx: Ctx) => number + readonly fail?: FailRule | undefined + } + | { + readonly kind: "external" + readonly host: string + readonly method: string + readonly path: string + readonly latency: Latency + readonly fail?: FailRule | undefined + } + | { + readonly kind: "internal" + readonly name: string + readonly self: Latency + readonly children: ReadonlyArray + } + | { readonly kind: "publish"; readonly topic: string; readonly consumer: Consumer } + +export interface Consumer { + readonly service: ServiceName + readonly topic: string + readonly self: Latency + readonly children: ReadonlyArray +} + +const fixed = + (median: number, spread = 0.35): Latency => + () => [median, spread] + +const pg = ( + operation: string, + collection: string, + statement: string, + latency: Latency = fixed(4, 0.5), + extra: Partial> = {}, +): Op => ({ + kind: "db", + system: "postgresql", + namespace: "shop", + operation, + collection, + statement, + latency, + ...extra, +}) + +const redis = (operation: string, statement: string): Op => ({ + kind: "db", + system: "redis", + // The map names database nodes by namespace; the DB index "0" reads as noise. + namespace: "cache", + operation, + statement, + latency: fixed(0.8, 0.4), +}) + +// ── failures ──────────────────────────────────────────────────────────────── + +/** + * Pool timeouts at a steady ~11% rather than independent coin flips, which at + * ~60 payment requests per 5-minute bucket swing from 0% to 25%. Error diffusion + * spaces failures evenly; a slow wave keeps the line from looking ruled. + */ +let poolDebt = 0 +const poolTimesOut = (ctx: Ctx): boolean => { + if (!ctx.incident) return false + const wave = 1 + 0.15 * Math.sin((2 * Math.PI * ctx.at) / (37 * MINUTE)) + poolDebt += 0.11 * wave * (0.5 + ctx.rng.next()) + if (poolDebt < 1) return false + poolDebt -= 1 + return true +} + +/** Every payment-svc query waits on the same exhausted pool. */ +const poolTimeout = + (logPrefix: string, callers: ReadonlyArray): FailRule => + (ctx) => + poolTimesOut(ctx) + ? { + type: "ConnectionTimeout", + message: "timed out acquiring a connection from the pool after 5000ms", + logPrefix, + durationMs: 5000 + ctx.rng.int(0, 40), + stacktrace: [ + "ConnectionTimeout: timed out acquiring a connection from the pool after 5000ms", + " at Pool.acquire (/app/node_modules/pg-pool/index.js:45:11)", + ...callers.map((frame) => ` at ${frame}`), + ].join("\n"), + } + : null + +const productCardTypeError: FailRule = (ctx) => + ctx.rng.chance(0.004) + ? { + type: "TypeError", + message: "Cannot read properties of undefined (reading 'price')", + stacktrace: [ + "TypeError: Cannot read properties of undefined (reading 'price')", + " at ProductCard (/app/.next/server/app/product/[slug]/page.js:1:4821)", + " at renderWithHooks (/app/node_modules/react-dom/cjs/react-dom-server.node.production.js:3:12099)", + " at renderElement (/app/node_modules/react-dom/cjs/react-dom-server.node.production.js:3:14980)", + ].join("\n"), + } + : null + +const variantKeyError: FailRule = (ctx) => + ctx.rng.chance(0.002) + ? { + type: "KeyError", + message: "'variant_id'", + stacktrace: [ + "Traceback (most recent call last):", + ' File "/app/catalog/api/products.py", line 74, in get_product', + " variant = variants_by_id[payload['variant_id']]", + "KeyError: 'variant_id'", + ].join("\n"), + } + : null + +const stockConflict: FailRule = (ctx) => + ctx.rng.chance(0.003) + ? { + type: "*inventory.ConflictError", + message: `stock reservation conflict: sku ${ctx.rng.int(1000, 9999)} version mismatch`, + stacktrace: [ + "goroutine 412 [running]:", + "github.com/acme/inventory-svc/internal/stock.(*Store).Reserve(0xc0001a2000, {0x10a3f80, 0xc000514120})", + "\t/src/internal/stock/store.go:118 +0x2c4", + "github.com/acme/inventory-svc/internal/http.(*Handler).reserve(0xc00012e0c0, {0x10a2c40, 0xc0002ae1c0}, 0xc000146300)", + "\t/src/internal/http/reserve.go:57 +0x1b8", + ].join("\n"), + } + : null + +const smtpDisconnect: FailRule = (ctx) => + ctx.rng.chance(0.006) + ? { + type: "smtplib.SMTPServerDisconnected", + message: "Connection unexpectedly closed", + stacktrace: [ + "Traceback (most recent call last):", + ' File "/app/notify/mailer.py", line 41, in send_receipt', + " smtp.send_message(message)", + ' File "/usr/local/lib/python3.12/smtplib.py", line 405, in getreply', + ' raise SMTPServerDisconnected("Connection unexpectedly closed")', + "smtplib.SMTPServerDisconnected: Connection unexpectedly closed", + ].join("\n"), + } + : null + +// ── call trees ────────────────────────────────────────────────────────────── + +const authSession: Op = { + kind: "server", + service: "auth-svc", + method: "GET", + route: "/session", + self: fixed(1.5), + children: [redis("GET", "GET session:{id}")], +} + +const productList: Op = { + kind: "server", + service: "catalog-api", + method: "GET", + route: "/products", + self: fixed(6, 0.4), + children: [ + redis("GET", "GET catalog:featured"), + pg( + "SELECT", + "products", + "SELECT id, slug, title, price_cents FROM products WHERE featured = $1 LIMIT $2", + ), + // The N+1 the flamegraph shows as a wall: one variants query per product. + pg("SELECT", "variants", "SELECT * FROM variants WHERE product_id = $1", fixed(1.6, 0.4), { + repeat: (ctx) => ctx.rng.int(12, 16), + }), + ], +} + +const productDetail: Op = { + kind: "server", + service: "catalog-api", + method: "GET", + route: "/products/{id}", + self: fixed(4, 0.4), + fail: variantKeyError, + children: [ + redis("GET", "GET product:{id}"), + pg("SELECT", "products", "SELECT * FROM products WHERE id = $1"), + pg("SELECT", "variants", "SELECT * FROM variants WHERE product_id = $1"), + ], +} + +const stockLookup: Op = { + kind: "server", + service: "inventory-svc", + method: "GET", + route: "/stock/{sku}", + self: fixed(1.2), + children: [pg("SELECT", "stock", "SELECT available, version FROM stock WHERE sku = $1")], +} + +const cartAdd: Op = { + kind: "server", + service: "cart-svc", + method: "POST", + route: "/carts/{id}/items", + self: fixed(2), + children: [stockLookup, redis("HSET", "HSET cart:{id} {sku} {qty}")], +} + +const cartGet: Op = { + kind: "server", + service: "cart-svc", + method: "GET", + route: "/carts/{id}", + self: fixed(1.4), + children: [redis("HGETALL", "HGETALL cart:{id}")], +} + +const reserveStock: Op = { + kind: "server", + service: "inventory-svc", + method: "POST", + route: "/reservations", + self: fixed(2), + fail: stockConflict, + children: [ + pg( + "UPDATE", + "stock", + "UPDATE stock SET available = available - $1, version = version + 1 WHERE sku = $2 AND version = $3", + fixed(6, 0.5), + ), + ], +} + +const charge: Op = { + kind: "server", + service: "payment-svc", + method: "POST", + route: "/charges", + self: fixed(3), + children: [ + pg( + "INSERT", + "payments", + "INSERT INTO payments (order_id, amount_cents, currency, status) VALUES ($1, $2, $3, 'pending')", + // Pool exhaustion: the insert now waits for a connection first. + (ctx) => (ctx.incident ? [620, 0.9] : [9, 0.45]), + { + fail: poolTimeout("charge failed", [ + "PaymentRepository.insert (/app/src/payments/repository.ts:88:24)", + "ChargeService.charge (/app/src/payments/charge.ts:142:18)", + "handleCharge (/app/src/routes/charge.ts:31:9)", + ]), + }, + ), + { + kind: "external", + host: "api.stripe.com", + method: "POST", + path: "/v1/payment_intents", + latency: fixed(240, 0.35), + }, + pg( + "UPDATE", + "payments", + "UPDATE payments SET status = $1, provider_ref = $2 WHERE id = $3", + fixed(5, 0.4), + ), + ], +} + +/** The order confirmation page polls this until the charge settles. */ +const chargeStatus: Op = { + kind: "server", + service: "payment-svc", + method: "GET", + route: "/charges/{id}", + self: fixed(1.5), + children: [ + pg( + "SELECT", + "payments", + "SELECT status, provider_ref FROM payments WHERE order_id = $1", + (ctx) => (ctx.incident ? [540, 0.9] : [3, 0.45]), + { + fail: poolTimeout("payment status lookup failed", [ + "PaymentRepository.findByOrder (/app/src/payments/repository.ts:41:24)", + "handleChargeStatus (/app/src/routes/charge-status.ts:18:9)", + ]), + }, + ), + ], +} + +const sendReceipt: Op = { + kind: "server", + service: "notification-svc", + method: "POST", + route: "/receipts", + self: fixed(4), + children: [ + { + kind: "external", + host: "smtp.postmarkapp.com", + method: "POST", + path: "/email", + latency: fixed(180, 0.4), + fail: smtpDisconnect, + }, + ], +} + +const orderCreated: Consumer = { + service: "order-worker", + topic: "orders.created", + self: fixed(3), + children: [ + pg( + "INSERT", + "orders", + "INSERT INTO orders (id, customer_id, total_cents, status) VALUES ($1, $2, $3, 'paid')", + ), + pg( + "INSERT", + "order_items", + "INSERT INTO order_items (order_id, sku, qty, price_cents) SELECT * FROM unnest($1, $2, $3, $4)", + ), + sendReceipt, + ], +} + +const checkout: Op = { + kind: "server", + service: "checkout-api", + method: "POST", + route: "/checkout", + self: fixed(4), + children: [ + authSession, + cartGet, + reserveStock, + charge, + { kind: "publish", topic: "orders.created", consumer: orderCreated }, + ], +} + +const storefront = (method: string, route: string, children: ReadonlyArray, fail?: FailRule): Op => ({ + kind: "server", + service: "storefront", + method, + route, + self: fixed(method === "GET" ? 14 : 5, 0.4), + children, + fail, +}) + +/** Entry points with their share of traffic. */ +export const ENTRY_POINTS: readonly [readonly [Op, number], ...(readonly [Op, number])[]] = [ + [storefront("GET", "/", [authSession, productList]), 24], + [ + storefront( + "GET", + "/product/[slug]", + [ + productDetail, + stockLookup, + { kind: "internal", name: "render ProductPage", self: fixed(9, 0.4), children: [] }, + ], + productCardTypeError, + ), + 30, + ], + [storefront("POST", "/api/cart", [authSession, cartAdd]), 14], + [storefront("GET", "/cart", [authSession, cartGet]), 8], + [storefront("POST", "/api/checkout", [checkout]), 14], + [storefront("GET", "/order/[id]", [authSession, chargeStatus]), 16], + [ + storefront("POST", "/api/login", [ + { + kind: "server", + service: "auth-svc", + method: "POST", + route: "/login", + self: fixed(3), + children: [ + pg("SELECT", "users", "SELECT id, password_hash FROM users WHERE email = $1"), + { kind: "internal", name: "argon2.verify", self: fixed(38, 0.15), children: [] }, + redis("SET", "SET session:{id} EX 86400"), + ], + }, + ]), + 6, + ], +] + +/** + * Traces per minute: peaks at 19:00 UTC (mid-afternoon US East, where the + * screenshots render) and bottoms out before dawn. UTC so the curve never + * depends on the seeding machine's timezone. + */ +export const traceRate = (at: number, peakPerMinute: number): number => { + const hour = new Date(at).getUTCHours() + new Date(at).getUTCMinutes() / 60 + const daily = 0.5 + 0.5 * Math.cos((2 * Math.PI * (hour - 19)) / 24) + return peakPerMinute * (0.3 + 0.7 * daily) +} diff --git a/scripts/seed-demo/telemetry.ts b/scripts/seed-demo/telemetry.ts new file mode 100644 index 0000000000..5e7e39205d --- /dev/null +++ b/scripts/seed-demo/telemetry.ts @@ -0,0 +1,866 @@ +/** + * Turns the scenario's call trees into OTLP spans, correlated logs and + * span-derived metrics, one window at a time so a 24h seed never holds the + * whole day in memory. + */ +import { Rng } from "./rng" +import { + DEPLOYS, + ENTRY_POINTS, + INCIDENT, + SERVICES, + traceRate, + type Consumer, + type Ctx, + type Failure, + type Op, + type ServiceDef, + type ServiceName, +} from "./scenario" +import { + attr, + deltaHistogram, + deltaSum, + gauge, + nano, + Severity, + SpanKind, + StatusCode, + type HistogramPoint, + type KeyValue, + type LogRecord, + type Metric, + type NumberPoint, + type SeverityText, + type Span, + type SpanEvent, +} from "./otlp" + +const MINUTE = 60_000 +const CLUSTER = "prod-us-east-1" +const NODES = ["ip-10-0-12-84", "ip-10-0-13-201", "ip-10-0-21-17", "ip-10-0-22-140"] as const +const LATENCY_BOUNDS_S = [0.005, 0.01, 0.025, 0.05, 0.075, 0.1, 0.25, 0.5, 0.75, 1, 2.5, 5, 7.5, 10] + +export interface WorldOptions { + readonly anchor: number + readonly windowMs: number + readonly incidentBeforeAnchorMs: number + readonly peakTracesPerMinute: number + readonly seed: number +} + +export interface Resource { + readonly key: string + readonly service: ServiceDef + readonly version: string + readonly sha: string + readonly pod: string + readonly attributes: ReadonlyArray +} + +type NonEmpty = readonly [T, ...T[]] + +const nonEmpty = (length: number, make: () => T): NonEmpty => { + const first = make() + return [first, ...Array.from({ length: length - 1 }, make)] +} + +interface Rollout { + readonly version: string + readonly sha: string + readonly at: number + readonly pods: NonEmpty +} + +/** Every service's rollouts over the window, with stable pod names and SHAs. */ +export class World { + readonly start: number + readonly incidentAt: number + private readonly rollouts = new Map>() + + constructor(readonly options: WorldOptions) { + this.start = options.anchor - options.windowMs + this.incidentAt = options.anchor - options.incidentBeforeAnchorMs + // Separate stream from traffic, so changing the traffic never renames a pod. + const naming = new Rng(options.seed ^ 0x5eed) + for (const service of SERVICES) { + const deploys = DEPLOYS.filter((deploy) => deploy.service === service.name) + .map((deploy) => ({ + version: deploy.version, + // The bad deploy follows --incident-minutes, so rollout, gauges and errors agree. + at: + options.anchor - + (deploy.service === INCIDENT.service && deploy.version === INCIDENT.version + ? options.incidentBeforeAnchorMs + : deploy.beforeAnchorMs), + })) + .sort((a, b) => a.at - b.at) + const rollout = (version: string, at: number): Rollout => { + const sha = naming.hex(20) + const templateHash = naming.hex(5).slice(0, 9) + const pods = nonEmpty(service.replicas, () => + makeResource( + service, + version, + sha, + `${service.name}-${templateHash}-${podSuffix(naming)}`, + naming.pick(NODES), + ), + ) + return { version, sha, at, pods } + } + this.rollouts.set(service.name, [ + rollout(service.baseVersion, Number.NEGATIVE_INFINITY), + ...deploys.map(({ version, at }) => rollout(version, at)), + ]) + } + } + + resource(service: ServiceName, at: number, rng: Rng): Resource { + // The constructor inserts every ServiceName. + const rollouts = this.rollouts.get(service)! + const live = rollouts.filter((rollout) => rollout.at <= at).at(-1) ?? rollouts[0] + return rng.pick(live.pods) + } + + allResources(): ReadonlyArray { + return [...this.rollouts.values()].flatMap((rollouts) => rollouts.flatMap((rollout) => rollout.pods)) + } +} + +const POD_ALPHABET = "bcdfghjklmnpqrstvwxz2456789" +const podSuffix = (rng: Rng) => + Array.from({ length: 5 }, () => POD_ALPHABET[Math.floor(rng.next() * POD_ALPHABET.length)]).join("") + +const makeResource = ( + service: ServiceDef, + version: string, + sha: string, + pod: string, + node: string, +): Resource => ({ + key: pod, + service, + version, + sha, + pod, + attributes: [ + attr("service.name", service.name), + attr("service.version", version), + attr("service.namespace", "shop"), + attr("service.instance.id", pod), + // Dual-emitted like every real Maple producer. + attr("deployment.environment.name", "production"), + attr("deployment.environment", "production"), + attr("vcs.ref.head.revision", sha), + attr("telemetry.sdk.name", "opentelemetry"), + attr("telemetry.sdk.language", service.language), + attr("cloud.provider", "aws"), + attr("cloud.region", "us-east-1"), + attr("k8s.cluster.name", CLUSTER), + attr("k8s.namespace.name", "shop"), + attr("k8s.deployment.name", service.name), + attr("k8s.pod.name", pod), + attr("k8s.node.name", node), + ], +}) + +// ── output ────────────────────────────────────────────────────────────────── + +export interface Grouped { + readonly resource: Resource + readonly items: T[] +} + +interface RequestAgg { + readonly resource: Resource + readonly minute: number + readonly method: string + readonly route: string + readonly status: number + readonly durationsMs: number[] +} + +export interface WindowOutput { + readonly spans: ReadonlyArray> + readonly logs: ReadonlyArray> + readonly metrics: ReadonlyArray> + readonly traceCount: number + readonly errorTraceCount: number +} + +class Sink { + readonly spans = new Map>() + readonly logs = new Map>() + readonly requests = new Map() + readonly orders = new Map() + + span(resource: Resource, span: Span) { + group(this.spans, resource).items.push(span) + } + + log(resource: Resource, record: LogRecord) { + group(this.logs, resource).items.push(record) + } + + request( + resource: Resource, + at: number, + method: string, + route: string, + status: number, + durationMs: number, + ) { + const minute = Math.floor(at / MINUTE) * MINUTE + const key = `${resource.key}|${minute}|${method}|${route}|${status}` + const agg = this.requests.get(key) ?? { resource, minute, method, route, status, durationsMs: [] } + agg.durationsMs.push(durationMs) + this.requests.set(key, agg) + } + + order(resource: Resource, at: number) { + const minute = Math.floor(at / MINUTE) * MINUTE + const key = `${resource.key}|${minute}` + const agg = this.orders.get(key) ?? { resource, minute, count: 0 } + agg.count += 1 + this.orders.set(key, agg) + } +} + +const group = (map: Map>, resource: Resource): Grouped => { + const existing = map.get(resource.key) + if (existing) return existing + const created: Grouped = { resource, items: [] } + map.set(resource.key, created) + return created +} + +// ── trace construction ────────────────────────────────────────────────────── + +interface Parent { + readonly spanId: string + readonly resource: Resource +} + +interface Failed { + readonly failure: Failure + readonly service: string +} + +interface Result { + readonly end: number + readonly failed: Failed | null +} + +const exceptionEvent = (at: number, failure: Failure): SpanEvent => ({ + timeUnixNano: nano(at), + name: "exception", + attributes: [ + attr("exception.type", failure.type), + attr("exception.message", failure.message), + attr("exception.stacktrace", failure.stacktrace), + ], +}) + +const fillRoute = (route: string, rng: Rng) => + route + .replace("{id}", String(rng.int(10_000, 99_999))) + .replace("{sku}", `SKU-${rng.int(1000, 9999)}`) + .replace("[id]", `ord_${rng.hex(4)}`) + .replace( + "[slug]", + rng.pick(["linen-overshirt", "trail-runner-2", "canvas-tote", "merino-beanie", "field-watch"]), + ) + +/** + * What a server or consumer span records when a call beneath it failed: the same + * exception if it was thrown in this service, otherwise its own UpstreamError + * naming the callee, the way a real handler wraps a failed downstream call. + */ +const recordedFailure = (service: string, failed: Failed | null): Failure | null => { + if (failed === null) return null + if (failed.service === service) return failed.failure + const message = `${failed.service} returned 500` + return { + type: "UpstreamError", + message, + stacktrace: [ + `UpstreamError: ${message}`, + " at callService (/app/src/lib/upstream.ts:41:13)", + " at async handle (/app/src/server/handler.ts:27:20)", + ].join("\n"), + } +} + +const consumerStatus = ( + service: string, + failed: Failed | null, + at: number, +): Pick => { + const recorded = recordedFailure(service, failed) + return recorded + ? { + events: [exceptionEvent(at, recorded)], + status: { code: StatusCode.error, message: recorded.message }, + } + : { status: { code: StatusCode.unset } } +} + +class TraceBuilder { + private readonly traceId: string + private readonly ctx: Ctx + + constructor( + private readonly world: World, + private readonly sink: Sink, + private readonly rng: Rng, + at: number, + ) { + this.traceId = rng.hex(16) + this.ctx = { at, incident: at >= world.incidentAt, rng } + } + + run(entry: Op): Result { + if (entry.kind !== "server") return { end: this.ctx.at, failed: null } + const resource = this.world.resource(entry.service, this.ctx.at, this.rng) + return this.server(entry, resource, undefined, this.ctx.at) + } + + private op(op: Op, parent: Parent, start: number): Result { + switch (op.kind) { + case "server": + return this.remote(op, parent, start) + case "db": + return this.db(op, parent, start) + case "external": + return this.external(op, parent, start) + case "internal": + return this.internal(op, parent, start) + case "publish": + return this.publish(op.topic, op.consumer, parent, start) + } + } + + private children(children: ReadonlyArray, parent: Parent, start: number): Result { + let cursor = start + for (const child of children) { + const result = this.op(child, parent, cursor + this.rng.next() * 0.4) + cursor = result.end + if (result.failed) return { end: cursor, failed: result.failed } + } + return { end: cursor, failed: null } + } + + private server( + op: Extract, + resource: Resource, + parentSpanId: string | undefined, + start: number, + ): Result { + const spanId = this.rng.hex(8) + const self = this.rng.latency(...op.self(this.ctx)) + const inner = this.children(op.children, { spanId, resource }, start + self * 0.4) + const own = inner.failed ? null : (op.fail?.(this.ctx) ?? null) + const end = inner.end + self * 0.6 + const failed: Failed | null = own ? { failure: own, service: op.service } : inner.failed + const recorded = own ?? recordedFailure(op.service, inner.failed) + const status = failed ? 500 : 200 + const path = fillRoute(op.route, this.rng) + + this.sink.span(resource, { + traceId: this.traceId, + spanId, + parentSpanId, + name: `${op.method} ${op.route}`, + kind: SpanKind.server, + startTimeUnixNano: nano(start), + endTimeUnixNano: nano(end), + attributes: [ + attr("http.request.method", op.method), + attr("http.route", op.route), + attr("url.path", path), + attr("url.scheme", "http"), + attr("http.response.status_code", status), + attr("server.address", `${op.service}.shop.svc.cluster.local`), + attr("network.protocol.version", "1.1"), + ...(recorded ? [attr("error.type", recorded.type)] : []), + ], + events: recorded ? [exceptionEvent(end, recorded)] : undefined, + status: recorded + ? { code: StatusCode.error, message: recorded.message } + : { code: StatusCode.unset }, + }) + this.sink.request(resource, start, op.method, op.route, status, end - start) + + const duration = Math.round(end - start) + const ids = { traceId: this.traceId, spanId } + // A failure bubbling up inside one service was already logged where it was thrown. + if (recorded && recorded !== inner.failed?.failure) { + this.log(resource, end, "ERROR", `${recorded.type}: ${recorded.message}`, ids, [ + attr("http.route", op.route), + attr("exception.type", recorded.type), + attr("exception.message", recorded.message), + ]) + } else if (failed) { + // Logged at the throw site. + } else if (duration > 2000) { + this.log(resource, end, "WARN", `slow request ${op.method} ${path} took ${duration}ms`, ids, [ + attr("http.route", op.route), + attr("duration_ms", duration), + ]) + } else if (this.rng.chance(0.55)) { + this.log(resource, end, "INFO", `${op.method} ${path} ${status} ${duration}ms`, ids, [ + attr("http.route", op.route), + attr("http.response.status_code", status), + attr("duration_ms", duration), + ]) + } + return { end, failed } + } + + private remote(op: Extract, parent: Parent, start: number): Result { + const spanId = this.rng.hex(8) + const target = this.world.resource(op.service, start, this.rng) + const served = this.server(op, target, spanId, start + this.networkHop()) + const end = served.end + this.networkHop() + const status = served.failed ? 500 : 200 + this.sink.span(parent.resource, { + traceId: this.traceId, + spanId, + parentSpanId: parent.spanId, + name: `${op.method} ${op.route}`, + kind: SpanKind.client, + startTimeUnixNano: nano(start), + endTimeUnixNano: nano(end), + attributes: [ + attr("http.request.method", op.method), + attr("url.template", op.route), + attr( + "url.full", + `http://${op.service}.shop.svc.cluster.local:8080${fillRoute(op.route, this.rng)}`, + ), + attr("server.address", `${op.service}.shop.svc.cluster.local`), + attr("server.port", 8080), + attr("peer.service", op.service), + attr("http.response.status_code", status), + ], + // The callee's server span and the caller's own UpstreamError carry the failure; + // erroring the client span too would add a third, exception-less issue per hop. + status: { code: StatusCode.unset }, + }) + return { end, failed: served.failed ? { failure: served.failed.failure, service: op.service } : null } + } + + /** At least 1ms, so a child still fits its parent once the UI truncates timestamps to ms. */ + private networkHop(): number { + return 1 + this.rng.latency(0.6, 0.4) + } + + private db(op: Extract, parent: Parent, start: number): Result { + const repeat = op.repeat?.(this.ctx) ?? 1 + const isPostgres = op.system === "postgresql" + let cursor = start + for (let i = 0; i < repeat; i++) { + const failure = op.fail?.(this.ctx) ?? null + const begin = cursor + this.rng.next() * 0.2 + const end = begin + (failure?.durationMs ?? this.rng.latency(...op.latency(this.ctx))) + const spanId = this.rng.hex(8) + this.sink.span(parent.resource, { + traceId: this.traceId, + spanId, + parentSpanId: parent.spanId, + name: isPostgres ? `${op.operation} ${op.namespace}.${op.collection}` : op.operation, + kind: SpanKind.client, + startTimeUnixNano: nano(begin), + endTimeUnixNano: nano(end), + attributes: [ + attr("db.system.name", op.system), + attr("db.system", op.system), + attr("db.namespace", op.namespace), + attr("db.operation.name", op.operation), + ...(op.collection ? [attr("db.collection.name", op.collection)] : []), + attr("db.query.text", op.statement), + attr("db.statement", op.statement), + attr( + "server.address", + isPostgres ? "orders-db.shop.svc.cluster.local" : "redis.shop.svc.cluster.local", + ), + attr("server.port", isPostgres ? 5432 : 6379), + ...(failure ? [attr("error.type", failure.type)] : []), + ], + events: failure ? [exceptionEvent(end, failure)] : undefined, + status: failure + ? { code: StatusCode.error, message: failure.message } + : { code: StatusCode.unset }, + }) + cursor = end + if (failure) { + this.failureLog(parent.resource, end, failure, spanId) + return { end, failed: { failure, service: parent.resource.service.name } } + } + } + return { end: cursor, failed: null } + } + + private external(op: Extract, parent: Parent, start: number): Result { + const failure = op.fail?.(this.ctx) ?? null + const end = start + (failure?.durationMs ?? this.rng.latency(...op.latency(this.ctx))) + const spanId = this.rng.hex(8) + this.sink.span(parent.resource, { + traceId: this.traceId, + spanId, + parentSpanId: parent.spanId, + name: `${op.method} ${op.host}`, + kind: SpanKind.client, + startTimeUnixNano: nano(start), + endTimeUnixNano: nano(end), + attributes: [ + attr("http.request.method", op.method), + attr("url.full", `https://${op.host}${op.path}`), + attr("server.address", op.host), + attr("server.port", 443), + ...(failure ? [attr("error.type", failure.type)] : [attr("http.response.status_code", 200)]), + ], + events: failure ? [exceptionEvent(end, failure)] : undefined, + status: failure + ? { code: StatusCode.error, message: failure.message } + : { code: StatusCode.unset }, + }) + if (failure) { + this.failureLog(parent.resource, end, failure, spanId) + return { end, failed: { failure, service: parent.resource.service.name } } + } + return { end, failed: null } + } + + private internal(op: Extract, parent: Parent, start: number): Result { + const spanId = this.rng.hex(8) + const self = this.rng.latency(...op.self(this.ctx)) + const inner = this.children(op.children, { spanId, resource: parent.resource }, start + self * 0.5) + const end = inner.end + self * 0.5 + this.sink.span(parent.resource, { + traceId: this.traceId, + spanId, + parentSpanId: parent.spanId, + name: op.name, + kind: SpanKind.internal, + startTimeUnixNano: nano(start), + endTimeUnixNano: nano(end), + attributes: [], + status: inner.failed ? { code: StatusCode.error } : { code: StatusCode.unset }, + }) + return { end, failed: inner.failed } + } + + /** The producer span ends the caller's work; the consumer runs later and never fails its caller. */ + private publish(topic: string, consumer: Consumer, parent: Parent, start: number): Result { + const spanId = this.rng.hex(8) + const end = start + this.rng.latency(1.8, 0.3) + this.sink.span(parent.resource, { + traceId: this.traceId, + spanId, + parentSpanId: parent.spanId, + name: `publish ${topic}`, + kind: SpanKind.producer, + startTimeUnixNano: nano(start), + endTimeUnixNano: nano(end), + attributes: [ + attr("messaging.system", "kafka"), + attr("messaging.destination.name", topic), + attr("messaging.operation.type", "send"), + attr("messaging.operation.name", "publish"), + attr("server.address", "kafka.shop.svc.cluster.local"), + ], + status: { code: StatusCode.unset }, + }) + + const begin = end + this.rng.latency(60, 0.8) + const resource = this.world.resource(consumer.service, begin, this.rng) + const consumerSpanId = this.rng.hex(8) + const self = this.rng.latency(...consumer.self(this.ctx)) + const inner = this.children( + consumer.children, + { spanId: consumerSpanId, resource }, + begin + self * 0.4, + ) + const consumed = inner.end + self * 0.6 + const partition = this.rng.int(0, 5) + this.sink.span(resource, { + traceId: this.traceId, + spanId: consumerSpanId, + parentSpanId: spanId, + name: `process ${consumer.topic}`, + kind: SpanKind.consumer, + startTimeUnixNano: nano(begin), + endTimeUnixNano: nano(consumed), + attributes: [ + attr("messaging.system", "kafka"), + attr("messaging.destination.name", consumer.topic), + attr("messaging.operation.type", "process"), + attr("messaging.consumer.group.name", consumer.service), + attr("messaging.destination.partition.id", String(partition)), + ], + ...consumerStatus(consumer.service, inner.failed, consumed), + }) + if (!inner.failed) { + this.sink.order(resource, begin) + const items = this.rng.int(1, 4) + this.log( + resource, + consumed, + "INFO", + `order ord_${this.rng.hex(4)} created: ${items} item${items > 1 ? "s" : ""}, $${this.rng.int(18, 240)}.${this.rng.int(10, 99)}`, + { traceId: this.traceId, spanId: consumerSpanId }, + [attr("messaging.destination.partition.id", String(partition))], + ) + } + return { end, failed: null } + } + + private failureLog(resource: Resource, at: number, failure: Failure, spanId: string) { + const message = failure.logPrefix + ? `${failure.logPrefix} for order ord_${this.rng.hex(4)}: ${failure.message}` + : `${failure.type}: ${failure.message}` + this.log(resource, at, "ERROR", message, { traceId: this.traceId, spanId }, [ + attr("exception.type", failure.type), + attr("exception.message", failure.message), + ]) + } + + private log( + resource: Resource, + at: number, + severity: SeverityText, + body: string, + ids: { traceId: string; spanId: string } | null, + attributes: ReadonlyArray, + ) { + this.sink.log(resource, logRecord(at, severity, body, ids, attributes)) + } +} + +const logRecord = ( + at: number, + severity: SeverityText, + body: string, + ids: { traceId: string; spanId: string } | null, + attributes: ReadonlyArray, +): LogRecord => ({ + timeUnixNano: nano(at), + observedTimeUnixNano: nano(at + 3), + severityNumber: Severity[severity], + severityText: severity, + body: { stringValue: body }, + attributes, + traceId: ids?.traceId, + spanId: ids?.spanId, +}) + +// ── windows ───────────────────────────────────────────────────────────────── + +/** Generate everything that happened in `[from, to)`. Call windows in order: the RNG is shared. */ +export const generateWindow = (world: World, rng: Rng, from: number, to: number): WindowOutput => { + const sink = new Sink() + let traceCount = 0 + let errorTraceCount = 0 + + for (let minute = from; minute < to; minute += MINUTE) { + const rate = traceRate(minute, world.options.peakTracesPerMinute) + const count = Math.max(0, Math.round(rate * (0.85 + 0.3 * rng.next()))) + const starts = Array.from({ length: count }, () => minute + rng.next() * MINUTE).sort((a, b) => a - b) + for (const at of starts) { + const result = new TraceBuilder(world, sink, rng, at).run(rng.weighted(ENTRY_POINTS)) + traceCount++ + if (result.failed) errorTraceCount++ + } + backgroundLogs(world, sink, rng, minute) + } + + return { + spans: [...sink.spans.values()], + logs: [...sink.logs.values()], + metrics: buildMetrics(world, sink, rng, from, to), + traceCount, + errorTraceCount, + } +} + +/** Log lines that belong to no request: pool pressure, consumer lag, sync jobs. */ +const backgroundLogs = (world: World, sink: Sink, rng: Rng, minute: number) => { + const at = (offset: number) => minute + offset + rng.next() * 900 + if (minute >= world.incidentAt) { + for (const offset of [0, 15_000, 30_000, 45_000]) { + const resource = world.resource(INCIDENT.service, minute, rng) + const waiting = rng.int(24, 48) + sink.log( + resource, + logRecord( + at(offset), + "WARN", + `pool exhausted: ${INCIDENT.poolSize}/${INCIDENT.poolSize} connections in use, ${waiting} waiting`, + null, + [ + attr("db.client.connection.pool.name", "payments"), + attr("db.client.connection.pending_requests", waiting), + ], + ), + ) + } + } + const worker = world.resource("order-worker", minute, rng) + const partition = rng.int(0, 5) + sink.log( + worker, + logRecord(at(0), "DEBUG", `consumer lag orders.created[${partition}]: ${rng.int(0, 14)}`, null, [ + attr("messaging.destination.partition.id", String(partition)), + ]), + ) + if (Math.floor(minute / MINUTE) % 5 === 0) { + const inventory = world.resource("inventory-svc", minute, rng) + sink.log( + inventory, + logRecord( + at(20_000), + "INFO", + `stock sync complete: ${rng.int(180, 420)} skus updated in ${rng.int(300, 900)}ms`, + null, + [], + ), + ) + } +} + +// ── metrics ───────────────────────────────────────────────────────────────── + +const histogramPoint = ( + minute: number, + durationsMs: ReadonlyArray, + attributes: ReadonlyArray, +): HistogramPoint => { + const seconds = durationsMs.map((ms) => ms / 1000) + const counts = LATENCY_BOUNDS_S.map(() => 0).concat(0) + for (const value of seconds) { + const bucket = LATENCY_BOUNDS_S.findIndex((bound) => value <= bound) + const index = bucket === -1 ? LATENCY_BOUNDS_S.length : bucket + counts[index] = (counts[index] ?? 0) + 1 + } + return { + startTimeUnixNano: nano(minute), + timeUnixNano: nano(minute + MINUTE), + attributes, + count: String(seconds.length), + sum: seconds.reduce((total, value) => total + value, 0), + min: Math.min(...seconds), + max: Math.max(...seconds), + bucketCounts: counts.map(String), + explicitBounds: LATENCY_BOUNDS_S, + } +} + +const buildMetrics = ( + world: World, + sink: Sink, + rng: Rng, + from: number, + to: number, +): ReadonlyArray> => { + const byResource = new Map() + const add = (resource: Resource, metric: Metric) => { + const entry = byResource.get(resource.key) ?? { resource, metrics: [] } + entry.metrics.push(metric) + byResource.set(resource.key, entry) + } + + // Request duration, straight from the server spans, so charts and traces agree. + const histograms = new Map() + const requestsPerMinute = new Map() + for (const agg of sink.requests.values()) { + const entry = histograms.get(agg.resource.key) ?? { resource: agg.resource, points: [] } + entry.points.push( + histogramPoint(agg.minute, agg.durationsMs, [ + attr("http.request.method", agg.method), + attr("http.route", agg.route), + attr("http.response.status_code", agg.status), + ]), + ) + histograms.set(agg.resource.key, entry) + const load = `${agg.resource.key}|${agg.minute}` + requestsPerMinute.set(load, (requestsPerMinute.get(load) ?? 0) + agg.durationsMs.length) + } + for (const { resource, points } of histograms.values()) { + add(resource, deltaHistogram("http.server.request.duration", "s", points)) + } + + // Process gauges for every pod that was live, scaled by the load it served. + for (const resource of world.allResources()) { + const cpu: NumberPoint[] = [] + const memory: NumberPoint[] = [] + for (let minute = from; minute < to; minute += MINUTE) { + const load = requestsPerMinute.get(`${resource.key}|${minute}`) + if (load === undefined) continue + const time = nano(minute + MINUTE) + cpu.push({ + timeUnixNano: time, + attributes: [], + asDouble: Math.min(0.95, 0.04 + load * 0.012 + rng.next() * 0.03), + }) + const baseMb = + resource.service.language === "go" ? 48 : resource.service.language === "python" ? 160 : 210 + memory.push({ + timeUnixNano: time, + attributes: [], + asInt: Math.round((baseMb + load * 0.8 + rng.next() * 12) * 1_048_576), + }) + } + if (cpu.length === 0) continue + add(resource, gauge("process.cpu.utilization", "1", cpu)) + add(resource, gauge("process.memory.usage", "By", memory)) + } + + // The payment pool: a few connections in use until 3.5.0 pins it at the limit. + for (const resource of world.allResources().filter((r) => r.service.name === INCIDENT.service)) { + const points: NumberPoint[] = [] + const pending: NumberPoint[] = [] + for (let minute = from; minute < to; minute += MINUTE) { + const load = requestsPerMinute.get(`${resource.key}|${minute}`) + if (load === undefined) continue + const exhausted = resource.version === INCIDENT.version + const used = exhausted + ? INCIDENT.poolSize + : Math.min(INCIDENT.poolSize - 4, 2 + Math.round(load * 0.35 + rng.next() * 2)) + const time = nano(minute + MINUTE) + const pool = attr("db.client.connection.pool.name", "payments") + points.push({ + timeUnixNano: time, + attributes: [pool, attr("db.client.connection.state", "used")], + asInt: used, + }) + points.push({ + timeUnixNano: time, + attributes: [pool, attr("db.client.connection.state", "idle")], + asInt: INCIDENT.poolSize - used, + }) + pending.push({ timeUnixNano: time, attributes: [pool], asInt: exhausted ? rng.int(24, 48) : 0 }) + } + if (points.length === 0) continue + add(resource, gauge("db.client.connection.count", "{connection}", points)) + add(resource, gauge("db.client.connection.pending_requests", "{request}", pending)) + } + + // The business number dashboards put next to error rate. + const orders = new Map() + for (const agg of sink.orders.values()) { + const entry = orders.get(agg.resource.key) ?? { resource: agg.resource, points: [] } + entry.points.push({ + startTimeUnixNano: nano(agg.minute), + timeUnixNano: nano(agg.minute + MINUTE), + attributes: [], + asInt: agg.count, + }) + orders.set(agg.resource.key, entry) + } + for (const { resource, points } of orders.values()) + add(resource, deltaSum("shop.orders.created", "{order}", points)) + + return [...byResource.values()].map(({ resource, metrics }) => ({ resource, items: metrics })) +}