diff --git a/apps/ingest/alchemy.run.ts b/apps/ingest/alchemy.run.ts index 056dcc79f..bfa15a1f7 100644 --- a/apps/ingest/alchemy.run.ts +++ b/apps/ingest/alchemy.run.ts @@ -11,8 +11,6 @@ import type { MapleRegion } from "@maple/infra/aws" import { COLLECTOR_DNS_LABEL, COLLECTOR_OTLP_HTTP_PORT, - INGEST_EC2_INSTANCE_TYPE, - INGEST_EC2_TASK_SIZE, parseIngestFleets, pgUrlRequireSsl, resolveAwsRegion, @@ -21,6 +19,8 @@ import { resolveCollectorTaskSize, resolveIngestCidrBlock, resolveIngestDesiredCount, + resolveIngestEc2InstanceType, + resolveIngestEc2TaskSize, resolveIngestNamespaceName, resolveIngestScaling, resolveIngestSelfTraceSampleRatio, @@ -227,7 +227,8 @@ export const createMapleIngest = ({ stage, domains, region, dbRole }: CreateMapl const replayBlobs = yield* replayBlobWriterCredentials(stage, region) const taskSize = resolveIngestTaskSize(stage) const selfTraceSampleRatio = resolveIngestSelfTraceSampleRatio(stage) - const scaling = resolveIngestScaling(stage) + const scaling = resolveIngestScaling(stage, region) + const ec2TaskSize = resolveIngestEc2TaskSize(stage, region) const name = (base: string) => resolveAwsResourceName(base, stage, region) const fleets = parseIngestFleets((yield* optionalPlain("MAPLE_INGEST_FLEETS")).MAPLE_INGEST_FLEETS) @@ -366,7 +367,7 @@ export const createMapleIngest = ({ stage, domains, region, dbRole }: CreateMapl { launchTemplateName: name("ingest-ec2"), imageId, - instanceType: INGEST_EC2_INSTANCE_TYPE, + instanceType: resolveIngestEc2InstanceType(stage, region), securityGroupIds: [instanceSecurityGroup.groupId], instanceProfileName: instanceProfile.instanceProfileName, associatePublicIpAddress: true, @@ -379,7 +380,7 @@ export const createMapleIngest = ({ stage, domains, region, dbRole }: CreateMapl // redeploy from resetting it to `minSize`); the bounds leave room for // a rolling deploy to double the fleet while old and new tasks // overlap on separate hosts. - const maxTasks = scaling?.max ?? resolveIngestDesiredCount(stage) + const maxTasks = scaling?.max ?? resolveIngestDesiredCount(stage, region) const autoScalingGroup = yield* AWS.AutoScaling.AutoScalingGroup("ingest-ec2-asg", { autoScalingGroupName: name("ingest-ec2"), launchTemplate, @@ -543,7 +544,7 @@ export const createMapleIngest = ({ stage, domains, region, dbRole }: CreateMapl tags: { Service: "maple-ingest", Region: region }, }) - const collectorTaskSize = resolveCollectorTaskSize(stage) + const collectorTaskSize = resolveCollectorTaskSize(stage, region) return yield* AWS.ECS.Service("otel-collector", { cluster, serviceName: name("otel-collector"), @@ -748,7 +749,7 @@ export const createMapleIngest = ({ stage, domains, region, dbRole }: CreateMapl // correct. runtimePlatform: { cpuArchitecture: "ARM64", operatingSystemFamily: "LINUX" } as const, - desiredCount: resolveIngestDesiredCount(stage), + desiredCount: resolveIngestDesiredCount(stage, region), // prd autoscales on CPU between this count and a burst ceiling; alchemy // stops pinning desiredCount while `scaling` is set, so the autoscaler's // decisions survive redeploys. Other stages stay fixed. @@ -929,6 +930,8 @@ export const createMapleIngest = ({ stage, domains, region, dbRole }: CreateMapl // own group (`ingest-ec2-sg`) is what admits the ALB. A rolling deploy // cannot start the new task beside the old one (the port is taken), so // managed scaling brings up a fresh host for it and drains the old one. + // That holds for EU prd's single host too: ECS's default 100/200 keeps the + // old task serving until the new host's task is healthy, so no min-healthy 0. const ec2Service = ec2Capacity ? yield* AWS.ECS.Service("ingest-ec2", { ...gateway, @@ -939,8 +942,8 @@ export const createMapleIngest = ({ stage, domains, region, dbRole }: CreateMapl { capacityProvider: ec2Capacity.capacityProvider.name, weight: 1 }, ], placementConstraints: [{ type: "distinctInstance" }], - cpu: INGEST_EC2_TASK_SIZE.cpu, - memory: INGEST_EC2_TASK_SIZE.memory, + cpu: ec2TaskSize.cpu, + memory: ec2TaskSize.memory, volumes: [{ name: "wal", host: { sourcePath: WAL_HOST_DIR } }], container: { stopTimeout, diff --git a/docs/infra.md b/docs/infra.md index 0f26fc80b..369a44f2c 100644 --- a/docs/infra.md +++ b/docs/infra.md @@ -649,6 +649,12 @@ architecture. preview stack (`scripts/ingest-preview.run.ts` and `deploy-pr-ingest.yml` are deleted). Two alchemy stacks claiming the same `maple-ingest-pr-` physical names is how orphan fleets accumulate. +- **The EU ingest fleet is sized to EU traffic** (`isEuPrd` in `packages/infra/src/aws/stage.ts`). + Over 30 days the EU ALB served 34k requests (~2 GB) against the US's 136M (~2.3 TB in), + yet ran the US footprint at ~$300/mo. EU prd runs one c7gd.medium (autoscaling 1-3, same + CPU target) and the collector at the non-prd size; Electric keeps the prd size. One host + still rolls a deploy: managed scaling adds a host for the new task, as in the US. To scale + it with traffic, raise the EU branches there or drop `isEuPrd` to inherit US sizing. - **Graviton (ARM64).** The gateway and collector tasks run `cpuArchitecture: "ARM64"` (`runtimePlatform` in `apps/ingest/alchemy.run.ts`), ~20% cheaper than x86_64. The blocker used to be the builder: cross-compiling Rust under QEMU is 10-30 min a build. diff --git a/packages/infra/src/aws/stage.test.ts b/packages/infra/src/aws/stage.test.ts index f60b32114..ef39a4957 100644 --- a/packages/infra/src/aws/stage.test.ts +++ b/packages/infra/src/aws/stage.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it } from "vitest" import { parseMapleStage } from "../cloudflare/stage.ts" import { + type MapleRegion, parseIngestFleets, parseMapleRegion, resolveAwsRegion, @@ -10,6 +11,8 @@ import { resolveElectricDbPoolSize, resolveIngestCidrBlock, resolveIngestDesiredCount, + resolveIngestEc2InstanceType, + resolveIngestEc2TaskSize, resolveIngestNamespaceName, resolveIngestScaling, resolveIngestSelfTraceSampleRatio, @@ -108,25 +111,63 @@ describe("collector service discovery", () => { }) it("sizes the collector task with 1 GiB everywhere so the memory limiter can fire", () => { - expect(resolveCollectorTaskSize(parseMapleStage("prd"))).toEqual({ cpu: 512, memory: 1024 }) - expect(resolveCollectorTaskSize(parseMapleStage("pr-12"))).toEqual({ cpu: 256, memory: 1024 }) - expect(resolveCollectorTaskSize(parseMapleStage("dev-alice"))).toEqual({ cpu: 256, memory: 1024 }) + expect(resolveCollectorTaskSize(parseMapleStage("prd"), "us")).toEqual({ cpu: 512, memory: 1024 }) + expect(resolveCollectorTaskSize(parseMapleStage("prd"), "eu")).toEqual({ cpu: 256, memory: 1024 }) + expect(resolveCollectorTaskSize(parseMapleStage("pr-12"), "us")).toEqual({ cpu: 256, memory: 1024 }) + expect(resolveCollectorTaskSize(parseMapleStage("dev-alice"), "us")).toEqual({ + cpu: 256, + memory: 1024, + }) }) }) describe("resolveIngestScaling", () => { it("autoscales production between the fixed count and a burst ceiling", () => { - const scaling = resolveIngestScaling(parseMapleStage("prd")) - expect(scaling).toBeDefined() - expect(scaling!.min).toBe(resolveIngestDesiredCount(parseMapleStage("prd"))) - expect(scaling!.max).toBeGreaterThan(scaling!.min) - expect(scaling!.cpuUtilization).toBeGreaterThan(0) - expect(scaling!.cpuUtilization).toBeLessThan(100) + const regions: ReadonlyArray = ["us", "eu"] + for (const region of regions) { + const scaling = resolveIngestScaling(parseMapleStage("prd"), region) + expect(scaling).toBeDefined() + expect(scaling!.min).toBe(resolveIngestDesiredCount(parseMapleStage("prd"), region)) + expect(scaling!.max).toBeGreaterThan(scaling!.min) + expect(scaling!.cpuUtilization).toBeGreaterThan(0) + expect(scaling!.cpuUtilization).toBeLessThan(100) + } + }) + + it("keeps US prd at 2-6 and sizes EU prd to its traffic at 1-3", () => { + const policy = { cpuUtilization: 60, scaleInCooldown: "5 minutes", scaleOutCooldown: "60 seconds" } + expect(resolveIngestScaling(parseMapleStage("prd"), "us")).toEqual({ min: 2, max: 6, ...policy }) + expect(resolveIngestScaling(parseMapleStage("prd"), "eu")).toEqual({ min: 1, max: 3, ...policy }) + expect(resolveIngestDesiredCount(parseMapleStage("prd"), "us")).toBe(2) + expect(resolveIngestDesiredCount(parseMapleStage("prd"), "eu")).toBe(1) }) it("keeps every other stage at a fixed count", () => { - expect(resolveIngestScaling(parseMapleStage("pr-12"))).toBeUndefined() - expect(resolveIngestScaling(parseMapleStage("dev-alice"))).toBeUndefined() + expect(resolveIngestScaling(parseMapleStage("pr-12"), "us")).toBeUndefined() + expect(resolveIngestScaling(parseMapleStage("dev-alice"), "us")).toBeUndefined() + expect(resolveIngestScaling(parseMapleStage("dev-alice"), "eu")).toBeUndefined() + expect(resolveIngestDesiredCount(parseMapleStage("pr-12"), "us")).toBe(1) + }) +}) + +describe("EC2 fleet sizing", () => { + it("runs US prd and previews on c7gd.large with the full-host task", () => { + for (const stage of ["prd", "pr-12"]) { + expect(resolveIngestEc2InstanceType(parseMapleStage(stage), "us")).toBe("c7gd.large") + expect(resolveIngestEc2TaskSize(parseMapleStage(stage), "us")).toEqual({ + cpu: 2048, + memory: 3072, + }) + } + }) + + it("runs EU prd on a c7gd.medium with a task that fits its registered memory", () => { + expect(resolveIngestEc2InstanceType(parseMapleStage("prd"), "eu")).toBe("c7gd.medium") + expect(resolveIngestEc2TaskSize(parseMapleStage("prd"), "eu")).toEqual({ cpu: 1024, memory: 1536 }) + }) + + it("keeps a non-prd EU stage on the default sizing", () => { + expect(resolveIngestEc2InstanceType(parseMapleStage("dev-alice"), "eu")).toBe("c7gd.large") }) }) diff --git a/packages/infra/src/aws/stage.ts b/packages/infra/src/aws/stage.ts index 164dbf14b..7917718d9 100644 --- a/packages/infra/src/aws/stage.ts +++ b/packages/infra/src/aws/stage.ts @@ -72,15 +72,27 @@ export function resolveAwsResourceName( } } +/** + * Whether this is the EU prd instance, whose ingest footprint is sized to its + * own traffic (34k requests in 30 days at launch vs the US's 136M) rather than + * the US fleet's. To scale it with traffic, raise the EU branches below (task + * count, scaling ceiling, instance type) or drop this check to inherit US sizing. + */ +function isEuPrd(stage: MapleStage, region: MapleRegion): boolean { + return stage.kind === "prd" && region === "eu" +} + /** * Desired ECS task count per stage. * - * prd runs 2 for availability across AZs; everything else runs 1. Note the - * per-org replay byte budget in the gateway is process-local - * (`apps/ingest/src/main.rs`), so the effective ceiling is roughly N x the - * configured limit — raising this raises that ceiling too. + * US prd runs 2 for availability across AZs; EU prd runs 1 (`isEuPrd`), and a + * dead task's WAL segments are claimed from S3 by its replacement. Everything + * else runs 1. Note the per-org replay byte budget in the gateway is + * process-local (`apps/ingest/src/main.rs`), so the effective ceiling is + * roughly N x the configured limit, so raising this raises that ceiling too. */ -export function resolveIngestDesiredCount(stage: MapleStage): number { +export function resolveIngestDesiredCount(stage: MapleStage, region: MapleRegion): number { + if (isEuPrd(stage, region)) return 1 return stage.kind === "prd" ? 2 : 1 } @@ -88,8 +100,9 @@ export function resolveIngestDesiredCount(stage: MapleStage): number { * Target-tracking autoscaling for the ingest service, or `undefined` for a * fixed desired count. * - * prd: 2–6 tasks on 60% average CPU. The floor is today's fixed count (AZ + * US prd: 2–6 tasks on 60% average CPU. The floor is today's fixed count (AZ * redundancy); the ceiling is ~6x current traffic at ~1,300 req/s per vCPU. + * EU prd: 1-3 tasks on the same target and cooldowns (`isEuPrd`). * Scale-out is eager (1 min) because a burst that outruns the gateway turns * into 5xx at the edge; scale-in is lazy (5 min) so a lull does not thrash. * Note the per-org replay byte budget is process-local, so the effective @@ -106,10 +119,10 @@ export interface IngestScaling { scaleOutCooldown: `${number} ${"seconds" | "minutes"}` } -export function resolveIngestScaling(stage: MapleStage): IngestScaling | undefined { - return stage.kind === "prd" - ? { min: 2, max: 6, cpuUtilization: 60, scaleInCooldown: "5 minutes", scaleOutCooldown: "60 seconds" } - : undefined +export function resolveIngestScaling(stage: MapleStage, region: MapleRegion): IngestScaling | undefined { + if (stage.kind !== "prd") return undefined + const [min, max] = isEuPrd(stage, region) ? [1, 3] : [2, 6] + return { min, max, cpuUtilization: 60, scaleInCooldown: "5 minutes", scaleOutCooldown: "60 seconds" } } export interface IngestTaskSize { @@ -171,19 +184,28 @@ export function parseIngestFleets(value: string | undefined): IngestFleets { } /** - * EC2 instance type for the gateway: Graviton3 with a 118 GB local NVMe - * instance store, which holds the WAL. The `d` is the point: the WAL fsyncs - * every frame, and instance-store fsync is tens of microseconds where Fargate's + * EC2 instance type for the gateway: Graviton3 with a local NVMe instance + * store, which holds the WAL. The `d` is the point: the WAL fsyncs every frame, + * and instance-store fsync is tens of microseconds where Fargate's * network-backed ephemeral storage is milliseconds. + * + * c7gd.large (2 vCPU, 4 GiB, 118 GB NVMe) everywhere but EU prd, which runs a + * c7gd.medium (1 vCPU, 2 GiB, 59 GB NVMe, see `isEuPrd`). The 48 GiB WAL cap + * (`WAL_MAX_BYTES` in `apps/ingest/alchemy.run.ts`) still fits the medium's disk. */ -export const INGEST_EC2_INSTANCE_TYPE = "c7gd.large" +export function resolveIngestEc2InstanceType(stage: MapleStage, region: MapleRegion): string { + return isEuPrd(stage, region) ? "c7gd.medium" : "c7gd.large" +} /** - * Task size on the EC2 fleet, every stage. One task per instance (host - * networking binds the port), so it claims the c7gd.large's 2 vCPU and most of - * its ~3.7 GiB registered memory, leaving room for a per-host monitoring daemon. + * Task size on the EC2 fleet. One task per instance (host networking binds the + * port), so it claims the host's vCPU and most of its registered memory + * (~3.7 GiB on a c7gd.large, ~1.8 GiB on a medium), leaving room for a + * per-host monitoring daemon. Must change with `resolveIngestEc2InstanceType`. */ -export const INGEST_EC2_TASK_SIZE: IngestTaskSize = { cpu: 2048, memory: 3072 } +export function resolveIngestEc2TaskSize(stage: MapleStage, region: MapleRegion): IngestTaskSize { + return isEuPrd(stage, region) ? { cpu: 1024, memory: 1536 } : { cpu: 2048, memory: 3072 } +} /** * Whether a stage gets an AWS ingest deployment at all. @@ -285,8 +307,11 @@ export function stageEnablesReplayBlobs(stage: MapleStage): boolean { * everywhere because the config's `memory_limiter` (768 MiB hard, 192 MiB * spike) is sized for it — a limit above task memory never fires. */ -export function resolveCollectorTaskSize(stage: MapleStage): IngestTaskSize { - return stage.kind === "prd" ? { cpu: 512, memory: 1024 } : { cpu: 256, memory: 1024 } +export function resolveCollectorTaskSize(stage: MapleStage, region: MapleRegion): IngestTaskSize { + // EU prd's self-telemetry is tiny, so it takes the non-prd size (`isEuPrd`). + return stage.kind === "prd" && !isEuPrd(stage, region) + ? { cpu: 512, memory: 1024 } + : { cpu: 256, memory: 1024 } } /** @@ -309,7 +334,8 @@ export function stageDeploysElectric(stage: MapleStage): boolean { * * Eight low-write control-plane tables, so this is sized for the BEAM's floor * rather than for throughput. Raise it when a shape's snapshot query, not its - * change stream, becomes the cost. + * change stream, becomes the cost. EU prd keeps the prd size: it serves the web + * app's sync, and 512 MiB is unmeasured against the BEAM's footprint under load. */ export function resolveElectricTaskSize(stage: MapleStage): IngestTaskSize { return stage.kind === "prd" ? { cpu: 512, memory: 1024 } : { cpu: 256, memory: 512 }