Skip to content

Commit 2fb4210

Browse files
matt-aitkenTrigger.dev RepoOps
authored andcommitted
feat(webapp): record API rate limit usage per environment in the metrics table
API rate limit usage is now recorded per environment in the `metrics` table as `api.rate_limit.allowed`, `api.rate_limit.denied`, `api.rate_limit.remaining_min`, `api.rate_limit.limit.per_second` and `api.rate_limit.limit.burst`, so you can chart requests against your limit and 429s over time on the Query page and in dashboards. Counts are aggregated in the rate-limit middleware and written as ClickHouse async inserts, so recording adds no per-request I/O. Off by default; enable with `API_RATE_LIMIT_METRICS_ENABLED=1`, or `allowlist` to record only organizations opted in through the `apiRateLimitMetricsEnabled` feature flag. Mono-RevId: 5cc33e4c426f673c694938cfd4eda686b591a7d0
1 parent 9d7a60b commit 2fb4210

19 files changed

Lines changed: 1351 additions & 16 deletions
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: feature
4+
---
5+
6+
API rate limit usage is now recorded per environment in the `metrics` table as `api.rate_limit.allowed`, `api.rate_limit.denied`, `api.rate_limit.remaining_min` and the limit itself (`api.rate_limit.limit.per_second` and `api.rate_limit.limit.burst`), so you can chart requests against your limit and 429s over time on the Query page and dashboards.

‎apps/webapp/app/entry.server.tsx‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import { initMollifierStaleSweepWorker } from "~/v3/mollifierStaleSweepWorker.se
1212
import { initBillingLimitWorker } from "~/v3/billingLimitWorker.server";
1313
import { initLogsSearchProjectorWorker } from "~/v3/logsSearchProjectorWorker.server";
1414
import { initQueueMetricsConsumer, initQueueMetricsEmitter } from "~/v3/queueMetrics.server";
15+
import { initApiRateLimitMetrics } from "~/services/apiRateLimitMetrics.server";
1516
import { bootstrap } from "./bootstrap";
1617
import { LocaleContextProvider } from "./components/primitives/LocaleProvider";
1718
import type { OperatingSystemPlatform } from "./components/primitives/OperatingSystemProvider";
@@ -281,6 +282,7 @@ initBillingLimitWorker();
281282
initLogsSearchProjectorWorker();
282283
initQueueMetricsEmitter();
283284
initQueueMetricsConsumer();
285+
initApiRateLimitMetrics();
284286

285287
bootstrap().catch((error) => {
286288
logError(error);

‎apps/webapp/app/env.server.ts‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -773,6 +773,25 @@ const EnvironmentSchema = z
773773
API_RATE_LIMIT_REJECTION_LOGS_ENABLED: z.string().default("1"),
774774
API_RATE_LIMIT_LIMITER_LOGS_ENABLED: z.string().default("0"),
775775

776+
API_RATE_LIMIT_METRICS_ENABLED: z.enum(["0", "1", "allowlist"]).default("0"),
777+
API_RATE_LIMIT_METRICS_BUCKET_SECONDS: z.coerce
778+
.number()
779+
.int()
780+
.positive()
781+
.multipleOf(10)
782+
.refine((seconds) => 60 % seconds === 0 || seconds % 60 === 0, {
783+
message: "must divide or be a multiple of 60 so buckets align to minute boundaries",
784+
})
785+
.default(10),
786+
API_RATE_LIMIT_METRICS_FLUSH_INTERVAL_MS: z.coerce.number().int().positive().default(10_000),
787+
API_RATE_LIMIT_METRICS_MAX_ENTRIES: z.coerce.number().int().positive().default(10_000),
788+
API_RATE_LIMIT_METRICS_WAIT_FOR_ASYNC_INSERT: z.string().default("0"),
789+
API_RATE_LIMIT_METRICS_INSERT_BUSY_TIMEOUT_MS: z.coerce
790+
.number()
791+
.int()
792+
.positive()
793+
.default(10_000),
794+
776795
API_RATE_LIMIT_JWT_WINDOW: z.string().default("1m"),
777796
API_RATE_LIMIT_JWT_TOKENS: z.coerce.number().int().default(60),
778797

‎apps/webapp/app/models/runtimeEnvironment.server.ts‎

Lines changed: 21 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -293,7 +293,10 @@ export async function findEnvironmentByApiKeyWithResolution(
293293

294294
export type PrivateApiKeyRateLimitScope = {
295295
environmentId: string;
296+
organizationId: string;
297+
projectId: string;
296298
apiRateLimiterConfig: unknown;
299+
featureFlags: unknown;
297300
};
298301

299302
export async function resolvePrivateApiKeyRateLimitScope(
@@ -313,8 +316,10 @@ export async function resolvePrivateApiKeyRateLimitScope(
313316
runtimeEnvironment: {
314317
select: {
315318
id: true,
319+
organizationId: true,
320+
projectId: true,
316321
project: { select: { deletedAt: true } },
317-
organization: { select: { apiRateLimiterConfig: true } },
322+
organization: { select: { apiRateLimiterConfig: true, featureFlags: true } },
318323
},
319324
},
320325
},
@@ -326,16 +331,21 @@ export async function resolvePrivateApiKeyRateLimitScope(
326331

327332
return {
328333
environmentId: match.runtimeEnvironment.id,
334+
organizationId: match.runtimeEnvironment.organizationId,
335+
projectId: match.runtimeEnvironment.projectId,
329336
apiRateLimiterConfig: match.runtimeEnvironment.organization.apiRateLimiterConfig,
337+
featureFlags: match.runtimeEnvironment.organization.featureFlags,
330338
};
331339
}
332340

333341
const environment = await tx.runtimeEnvironment.findFirst({
334342
where: { apiKey },
335343
select: {
336344
id: true,
345+
organizationId: true,
346+
projectId: true,
337347
project: { select: { deletedAt: true } },
338-
organization: { select: { apiRateLimiterConfig: true } },
348+
organization: { select: { apiRateLimiterConfig: true, featureFlags: true } },
339349
},
340350
});
341351

@@ -346,7 +356,10 @@ export async function resolvePrivateApiKeyRateLimitScope(
346356

347357
return {
348358
environmentId: environment.id,
359+
organizationId: environment.organizationId,
360+
projectId: environment.projectId,
349361
apiRateLimiterConfig: environment.organization.apiRateLimiterConfig,
362+
featureFlags: environment.organization.featureFlags,
350363
};
351364
}
352365

@@ -356,8 +369,10 @@ export async function resolvePrivateApiKeyRateLimitScope(
356369
runtimeEnvironment: {
357370
select: {
358371
id: true,
372+
organizationId: true,
373+
projectId: true,
359374
project: { select: { deletedAt: true } },
360-
organization: { select: { apiRateLimiterConfig: true } },
375+
organization: { select: { apiRateLimiterConfig: true, featureFlags: true } },
361376
},
362377
},
363378
},
@@ -370,7 +385,10 @@ export async function resolvePrivateApiKeyRateLimitScope(
370385

371386
return {
372387
environmentId: revokedEnvironment.id,
388+
organizationId: revokedEnvironment.organizationId,
389+
projectId: revokedEnvironment.projectId,
373390
apiRateLimiterConfig: revokedEnvironment.organization.apiRateLimiterConfig,
391+
featureFlags: revokedEnvironment.organization.featureFlags,
374392
};
375393
}
376394

‎apps/webapp/app/services/apiRateLimit.server.ts‎

Lines changed: 25 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,9 +3,14 @@ import { env } from "~/env.server";
33
import { resolvePrivateApiKeyRateLimitScope } from "~/models/runtimeEnvironment.server";
44
import { batchStreamGrants } from "~/runEngine/concerns/batchStreamGrantsInstance.server";
55
import { authenticateAuthorizationHeader } from "./apiAuth.server";
6-
import { authorizationRateLimitMiddleware } from "./authorizationRateLimitMiddleware.server";
6+
import { recordApiRateLimitObservation } from "./apiRateLimitMetrics.server";
7+
import {
8+
authorizationRateLimitMiddleware,
9+
type RateLimitTenant,
10+
} from "./authorizationRateLimitMiddleware.server";
711
import { deploymentApiPaths } from "./deploymentApiPaths.server";
812
import type { Duration } from "./rateLimiter.server";
13+
import { FEATURE_FLAG, FeatureFlagCatalog } from "~/v3/featureFlags";
914

1015
const BATCH_STREAM_ITEMS_PATH = /^\/api\/v3\/batches\/([^/]+)\/items$/;
1116

@@ -17,12 +22,23 @@ export function jwtActorRateLimitIdentifier(environmentId: string, actorSub: str
1722
return `jwt-actor:${environmentId}:${actorSub}`;
1823
}
1924

25+
/** The organization override for the metrics opt-in flag; anything unparseable reads as off. */
26+
export function readApiRateLimitMetricsFlag(featureFlags: unknown): boolean {
27+
if (!featureFlags || typeof featureFlags !== "object" || Array.isArray(featureFlags)) {
28+
return false;
29+
}
30+
const parsed = FeatureFlagCatalog[FEATURE_FLAG.apiRateLimitMetricsEnabled].safeParse(
31+
(featureFlags as Record<string, unknown>)[FEATURE_FLAG.apiRateLimitMetricsEnabled]
32+
);
33+
return parsed.success ? parsed.data : false;
34+
}
35+
2036
// The per-request bucket decision for the API limiter. Exported so the branch below
2137
// (a delegated JWT keys on env+acting-user, everything else keeps its prior key) is
2238
// testable without standing up the middleware and its Redis.
2339
export async function resolveApiRateLimitOverride(
2440
authorizationValue: string
25-
): Promise<{ config?: unknown; identifier?: string } | undefined> {
41+
): Promise<{ config?: unknown; identifier?: string; tenant?: RateLimitTenant } | undefined> {
2642
const rawApiKey = authorizationValue.replace(/^Bearer /, "");
2743

2844
if (rawApiKey.startsWith("tr_")) {
@@ -35,6 +51,12 @@ export async function resolveApiRateLimitOverride(
3551
return {
3652
config: scope.apiRateLimiterConfig,
3753
identifier: scope.environmentId,
54+
tenant: {
55+
organizationId: scope.organizationId,
56+
projectId: scope.projectId,
57+
environmentId: scope.environmentId,
58+
metricsEnabled: readApiRateLimitMetricsFlag(scope.featureFlags),
59+
},
3860
};
3961
}
4062

@@ -98,6 +120,7 @@ export const apiRateLimiter = authorizationRateLimitMiddleware({
98120
maxItems: 1000,
99121
},
100122
limiterConfigOverride: resolveApiRateLimitOverride,
123+
onResult: recordApiRateLimitObservation,
101124
pathMatchers: [/^\/api/],
102125
// Allow /api/v1/tasks/:id/callback/:secret
103126
pathWhiteList: [
Lines changed: 153 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,153 @@
1+
import type { MetricsV1Input } from "@internal/clickhouse";
2+
import { env } from "~/env.server";
3+
import { clickhouseFactory } from "~/services/clickhouse/clickhouseFactoryInstance.server";
4+
import { singleton } from "~/utils/singleton";
5+
import { meter } from "~/v3/tracer.server";
6+
import { ApiRateLimitMetricsAggregator } from "./apiRateLimitMetricsAggregator.server";
7+
import {
8+
apiRateLimitMetricsInsertSettings,
9+
exportApiRateLimitMetricRows,
10+
} from "./apiRateLimitMetricsExporter.server";
11+
import type {
12+
RateLimitObservation,
13+
RateLimitTenant,
14+
} from "./authorizationRateLimitMiddleware.server";
15+
import { logger } from "./logger.server";
16+
import { signalsEmitter } from "./signals.server";
17+
18+
function enabledByEnv(): boolean {
19+
return env.API_RATE_LIMIT_METRICS_ENABLED !== "0";
20+
}
21+
22+
export function recordApiRateLimitObservation(observation: RateLimitObservation): void {
23+
if (!enabledByEnv()) {
24+
return;
25+
}
26+
getAggregator().record(observation);
27+
}
28+
29+
/**
30+
* Builds the aggregator and starts its flush timer ahead of the first request.
31+
*/
32+
export function initApiRateLimitMetrics(): void {
33+
if (!enabledByEnv()) {
34+
return;
35+
}
36+
getAggregator();
37+
}
38+
39+
function getAggregator(): ApiRateLimitMetricsAggregator {
40+
return singleton("apiRateLimitMetricsAggregator", createAggregator);
41+
}
42+
43+
/**
44+
* With the env value "allowlist" only tenants whose organization carries the
45+
* apiRateLimitMetricsEnabled feature flag are recorded. The flag travels with the cached
46+
* rate-limit resolution, so this costs nothing per request. Turning recording off is a redeploy.
47+
*/
48+
function createRuntimeGate(): (tenant: RateLimitTenant) => boolean {
49+
if (env.API_RATE_LIMIT_METRICS_ENABLED === "1") {
50+
return () => true;
51+
}
52+
return (tenant) => tenant.metricsEnabled;
53+
}
54+
55+
function exportRows(rows: MetricsV1Input[], onInsertError: (rows: number) => void): Promise<void> {
56+
return exportApiRateLimitMetricRows(rows, {
57+
resolveClient: (organizationId) =>
58+
clickhouseFactory.getClickhouseForOrganizationSync(organizationId, "events"),
59+
settings: apiRateLimitMetricsInsertSettings({
60+
waitForAsyncInsert: env.API_RATE_LIMIT_METRICS_WAIT_FOR_ASYNC_INSERT === "1",
61+
busyTimeoutMs: env.API_RATE_LIMIT_METRICS_INSERT_BUSY_TIMEOUT_MS,
62+
}),
63+
onInsertError: (failedRows, error) => {
64+
onInsertError(failedRows);
65+
logger.error(
66+
"api rate limit metrics: clickhouse rejected the insert request, dropping rows",
67+
{
68+
rows: failedRows,
69+
error: error instanceof Error ? error.message : String(error),
70+
}
71+
);
72+
},
73+
});
74+
}
75+
76+
function createAggregator(): ApiRateLimitMetricsAggregator {
77+
const droppedCounter = meter.createCounter("api_rate_limit_metrics.dropped", {
78+
description:
79+
"API rate limit observations dropped, by reason: the in-process cap was reached, or shutdown came before data store routing was ready",
80+
});
81+
const flushedCounter = meter.createCounter("api_rate_limit_metrics.rows_flushed", {
82+
description: "API rate limit metric rows sent to ClickHouse as async inserts",
83+
});
84+
const insertFailedCounter = meter.createCounter("api_rate_limit_metrics.rows_insert_failed", {
85+
description:
86+
"API rate limit metric rows lost because ClickHouse rejected the insert request; failures while writing a queued async insert are not visible here",
87+
});
88+
89+
const aggregator = new ApiRateLimitMetricsAggregator({
90+
bucketSeconds: env.API_RATE_LIMIT_METRICS_BUCKET_SECONDS,
91+
maxEntries: env.API_RATE_LIMIT_METRICS_MAX_ENTRIES,
92+
isEnabled: createRuntimeGate(),
93+
sink: (rows) => exportRows(rows, (failed) => insertFailedCounter.add(failed)),
94+
onDropped: (count) => droppedCounter.add(count, { reason: "cap" }),
95+
});
96+
97+
let routingReady = false;
98+
clickhouseFactory
99+
.isReady()
100+
.then(() => {
101+
routingReady = true;
102+
})
103+
.catch((error) => {
104+
logger.error("api rate limit metrics: data store registry never became ready", {
105+
error: error instanceof Error ? error.message : String(error),
106+
});
107+
});
108+
109+
let deferredWarned = false;
110+
const flush = () => {
111+
if (!routingReady) {
112+
if (aggregator.size > 0 && !deferredWarned) {
113+
deferredWarned = true;
114+
logger.warn("api rate limit metrics: flush deferred, data store routing not ready", {
115+
pendingEntries: aggregator.size,
116+
});
117+
}
118+
return;
119+
}
120+
try {
121+
const flushed = aggregator.flush();
122+
if (flushed > 0) {
123+
flushedCounter.add(flushed);
124+
}
125+
} catch (error) {
126+
logger.error("api rate limit metrics: flush failed", {
127+
error: error instanceof Error ? error.message : String(error),
128+
});
129+
}
130+
};
131+
132+
const interval = setInterval(flush, env.API_RATE_LIMIT_METRICS_FLUSH_INTERVAL_MS);
133+
interval.unref();
134+
135+
const shutdown = () => {
136+
clearInterval(interval);
137+
if (!routingReady) {
138+
const dropped = aggregator.discard();
139+
if (dropped > 0) {
140+
droppedCounter.add(dropped, { reason: "shutdown_before_ready" });
141+
logger.warn("api rate limit metrics: shutting down before routing was ready, dropping", {
142+
droppedObservations: dropped,
143+
});
144+
}
145+
return;
146+
}
147+
flush();
148+
};
149+
signalsEmitter.on("SIGTERM", shutdown);
150+
signalsEmitter.on("SIGINT", shutdown);
151+
152+
return aggregator;
153+
}

0 commit comments

Comments
 (0)