Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 12 additions & 9 deletions apps/ingest/alchemy.run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -21,6 +19,8 @@ import {
resolveCollectorTaskSize,
resolveIngestCidrBlock,
resolveIngestDesiredCount,
resolveIngestEc2InstanceType,
resolveIngestEc2TaskSize,
resolveIngestNamespaceName,
resolveIngestScaling,
resolveIngestSelfTraceSampleRatio,
Expand Down Expand Up @@ -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)

Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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"),
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down
6 changes: 6 additions & 0 deletions docs/infra.md
Original file line number Diff line number Diff line change
Expand Up @@ -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-<n>` 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.
Expand Down
63 changes: 52 additions & 11 deletions packages/infra/src/aws/stage.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { describe, expect, it } from "vitest"
import { parseMapleStage } from "../cloudflare/stage.ts"
import {
type MapleRegion,
parseIngestFleets,
parseMapleRegion,
resolveAwsRegion,
Expand All @@ -10,6 +11,8 @@ import {
resolveElectricDbPoolSize,
resolveIngestCidrBlock,
resolveIngestDesiredCount,
resolveIngestEc2InstanceType,
resolveIngestEc2TaskSize,
resolveIngestNamespaceName,
resolveIngestScaling,
resolveIngestSelfTraceSampleRatio,
Expand Down Expand Up @@ -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<MapleRegion> = ["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")
})
})

Expand Down
68 changes: 47 additions & 21 deletions packages/infra/src/aws/stage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,24 +72,37 @@ 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
}

/**
* 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
Expand All @@ -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 {
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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 }
}

/**
Expand All @@ -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 }
Expand Down
Loading