From 559c1e7bb14765acb234bf004c8e839fadfc2922 Mon Sep 17 00:00:00 2001 From: Andy Bitz Date: Wed, 2 Sep 2026 13:10:03 +0200 Subject: [PATCH 1/3] Use the runtime-provided ingest transport when available --- .changeset/tidy-donuts-brush.md | 5 + .../vercel-flags-core/src/utils/ingest.ts | 13 ++- .../src/utils/runtime-ingest.ts | 27 +++++ .../vercel-flags-core/src/utils/scheduler.ts | 6 +- .../src/utils/usage-tracker.test.ts | 101 ++++++++++++++++++ .../src/utils/usage-tracker.ts | 19 +++- 6 files changed, 167 insertions(+), 4 deletions(-) create mode 100644 .changeset/tidy-donuts-brush.md create mode 100644 packages/vercel-flags-core/src/utils/runtime-ingest.ts diff --git a/.changeset/tidy-donuts-brush.md b/.changeset/tidy-donuts-brush.md new file mode 100644 index 00000000..396faf21 --- /dev/null +++ b/.changeset/tidy-donuts-brush.md @@ -0,0 +1,5 @@ +--- +'@vercel/flags-core': patch +--- + +Use the runtime-provided ingest transport when available diff --git a/packages/vercel-flags-core/src/utils/ingest.ts b/packages/vercel-flags-core/src/utils/ingest.ts index 20a769b5..c5f855ba 100644 --- a/packages/vercel-flags-core/src/utils/ingest.ts +++ b/packages/vercel-flags-core/src/utils/ingest.ts @@ -3,6 +3,7 @@ import { version } from '../../package.json'; import type { Auth } from '../controller/auth'; import type { MetricEnvironment } from '../types'; import { getRetryDelayMs } from './backoff'; +import { getRuntimeIngest } from './runtime-ingest'; import type { FlushReason } from './scheduler'; import type { IngestEvent, UsageEvent } from './usage/events'; @@ -76,7 +77,17 @@ export async function sendIngestEvents( flushId: number, flushReason: FlushReason, ): Promise { - const eventsToSend = events.map((event) => event.ingestEvent()); + let eventsToSend = events.map((event) => event.ingestEvent()); + + const runtimeIngest = getRuntimeIngest(); + if (runtimeIngest) { + const headers = await getIngestHeaders(options, flushReason); + // Events the runtime does not accept fall through to the HTTP transport. + eventsToSend = eventsToSend.filter( + (event) => !runtimeIngest({ headers, body: [event] }), + ); + if (eventsToSend.length === 0) return; + } for (let i = 0; i < eventsToSend.length; i += MAX_EVENTS_PER_REQUEST) { await sendIngestChunk( diff --git a/packages/vercel-flags-core/src/utils/runtime-ingest.ts b/packages/vercel-flags-core/src/utils/runtime-ingest.ts new file mode 100644 index 00000000..51c8c817 --- /dev/null +++ b/packages/vercel-flags-core/src/utils/runtime-ingest.ts @@ -0,0 +1,27 @@ +import type { IngestEvent } from './usage/events'; + +export type RuntimeIngest = (payload: { + headers: Record; + body: IngestEvent[]; +}) => boolean; + +const FLAGS_CONTEXT_SYMBOL = Symbol.for('@vercel/flags-context'); + +/** + * Returns the ingest transport provided by the runtime, if available. + */ +export function getRuntimeIngest(): RuntimeIngest | undefined { + try { + const context = ( + globalThis as typeof globalThis & { + [key: symbol]: { ingest?: unknown } | undefined; + } + )[FLAGS_CONTEXT_SYMBOL]; + + return typeof context?.ingest === 'function' + ? (context.ingest as RuntimeIngest) + : undefined; + } catch { + return undefined; + } +} diff --git a/packages/vercel-flags-core/src/utils/scheduler.ts b/packages/vercel-flags-core/src/utils/scheduler.ts index c7ce7594..6a740dd1 100644 --- a/packages/vercel-flags-core/src/utils/scheduler.ts +++ b/packages/vercel-flags-core/src/utils/scheduler.ts @@ -5,7 +5,11 @@ const IDLE_FLUSH_WAIT_MS = 5000; const IDLE_FLUSH_JITTER_RATIO = 0.2; const MAX_FLUSH_WAIT_MS = 60000; -export type FlushReason = 'idle_timeout' | 'max_timeout' | 'shutdown'; +export type FlushReason = + | 'idle_timeout' + | 'max_timeout' + | 'shutdown' + | 'immediate'; /** * Schedule helper that flushes when any of the following occur: diff --git a/packages/vercel-flags-core/src/utils/usage-tracker.test.ts b/packages/vercel-flags-core/src/utils/usage-tracker.test.ts index ca269d06..ea4b7448 100644 --- a/packages/vercel-flags-core/src/utils/usage-tracker.test.ts +++ b/packages/vercel-flags-core/src/utils/usage-tracker.test.ts @@ -1220,3 +1220,104 @@ describe('UsageTracker', () => { }); }); }); + +describe('runtime ingest transport', () => { + const FLAGS_CONTEXT_SYMBOL = Symbol.for('@vercel/flags-context'); + + type RuntimeIngestPayload = { + headers: Record; + body: { type: string; ts: number; payload: object }[]; + }; + + let ingestMock: ReturnType< + typeof vi.fn<(p: RuntimeIngestPayload) => boolean> + >; + + beforeEach(() => { + ingestMock = vi.fn<(p: RuntimeIngestPayload) => boolean>(); + Object.defineProperty(globalThis, FLAGS_CONTEXT_SYMBOL, { + value: { ingest: ingestMock }, + configurable: true, + }); + }); + + afterEach(() => { + delete (globalThis as Record)[FLAGS_CONTEXT_SYMBOL]; + }); + + it('delivers events through the runtime without fetch or waitUntil', async () => { + ingestMock.mockReturnValue(true); + + const tracker = createTracker(); + tracker.trackEvaluation({ + flagKey: 'my-flag', + variant: 'on', + reason: ResolutionReason.RULE_MATCH, + }); + + await vi.waitFor(() => expect(ingestMock).toHaveBeenCalledTimes(1)); + + const { headers, body } = ingestMock.mock.calls[0]![0]; + expect(headers.Authorization).toBe('Bearer test-key'); + expect(headers[FLUSH_REASON_HEADER]).toBe('immediate'); + expect(body).toHaveLength(1); + expect(body[0]!.type).toBe('FLAG_EVALUATION'); + + expect(fetchMock).not.toHaveBeenCalled(); + expect(waitUntilMock).not.toHaveBeenCalled(); + }); + + it('delivers each event separately', async () => { + ingestMock.mockReturnValue(true); + + const tracker = createTracker(); + tracker.trackRead(); + tracker.trackEvaluation({ + flagKey: 'my-flag', + variant: 'on', + reason: ResolutionReason.RULE_MATCH, + }); + + await vi.waitFor(() => expect(ingestMock).toHaveBeenCalledTimes(2)); + const types = ingestMock.mock.calls.flatMap((call) => + call[0].body.map((event) => event.type), + ); + expect(types).toEqual( + expect.arrayContaining(['FLAGS_CONFIG_READ', 'FLAG_EVALUATION']), + ); + expect(fetchMock).not.toHaveBeenCalled(); + }); + + it('falls back to fetch for events the runtime does not accept', async () => { + ingestMock.mockReturnValue(false); + fetchMock.mockImplementation(() => jsonResponse({ ok: true })); + + const tracker = createTracker(); + tracker.trackEvaluation({ + flagKey: 'my-flag', + variant: 'on', + reason: ResolutionReason.RULE_MATCH, + }); + + await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledTimes(1)); + const events = getBody() as SerializedEvaluationEvent[]; + expect(events).toHaveLength(1); + expect(events[0]!.type).toBe('FLAG_EVALUATION'); + }); + + it('uses the scheduler when the runtime does not provide a transport', async () => { + delete (globalThis as Record)[FLAGS_CONTEXT_SYMBOL]; + fetchMock.mockImplementation(() => jsonResponse({ ok: true })); + + const tracker = createTracker(); + tracker.trackEvaluation({ + flagKey: 'my-flag', + variant: 'on', + reason: ResolutionReason.RULE_MATCH, + }); + + expect(waitUntilMock).toHaveBeenCalledTimes(1); + await tracker.shutdown(); + expect(fetchMock).toHaveBeenCalledTimes(1); + }); +}); diff --git a/packages/vercel-flags-core/src/utils/usage-tracker.ts b/packages/vercel-flags-core/src/utils/usage-tracker.ts index ab0e78a0..e012efdb 100644 --- a/packages/vercel-flags-core/src/utils/usage-tracker.ts +++ b/packages/vercel-flags-core/src/utils/usage-tracker.ts @@ -1,5 +1,6 @@ import { type IngestOptions, sendIngestEvents } from './ingest'; import { getRequestContext } from './request-context'; +import { getRuntimeIngest } from './runtime-ingest'; import { type FlushReason, Scheduler } from './scheduler'; import { FlagsConfigReadEvent, @@ -60,7 +61,7 @@ export class UsageTracker { this.readEvents.push(new FlagsConfigReadEvent(headers, options)); - this.scheduler.scheduleFlush(); + this.requestFlush(); } catch (error) { // trackRead should never throw, but log the error console.error('@vercel/flags-core: Failed to record event:', error); @@ -90,7 +91,7 @@ export class UsageTracker { } // always schedule to reset the timer - this.scheduler.scheduleFlush(); + this.requestFlush(); } catch (error) { console.error( '@vercel/flags-core: Failed to record evaluation event:', @@ -99,6 +100,20 @@ export class UsageTracker { } } + /** + * Flushes immediately when the runtime provides an ingest transport, + * otherwise falls back to the time-based scheduler. + */ + private requestFlush(): void { + if (getRuntimeIngest()) { + void this.flushEvents('immediate').catch((error) => { + console.error('@vercel/flags-core: Failed to flush events:', error); + }); + } else { + this.scheduler.scheduleFlush(); + } + } + /** * Send all events to the ingest service */ From 9d0e81e72f41a457dc222ea663d888bede05e253 Mon Sep 17 00:00:00 2001 From: Andy Bitz Date: Tue, 8 Sep 2026 17:09:15 +0200 Subject: [PATCH 2/3] Send only the SDK key on the runtime ingest transport --- .../vercel-flags-core/src/utils/ingest.ts | 28 +++++++++++++++- .../src/utils/usage-tracker.test.ts | 32 +++++++++++++++++++ 2 files changed, 59 insertions(+), 1 deletion(-) diff --git a/packages/vercel-flags-core/src/utils/ingest.ts b/packages/vercel-flags-core/src/utils/ingest.ts index c5f855ba..755707e3 100644 --- a/packages/vercel-flags-core/src/utils/ingest.ts +++ b/packages/vercel-flags-core/src/utils/ingest.ts @@ -71,6 +71,32 @@ async function getIngestHeaders( }; } +/** + * Headers for the runtime-provided ingest transport. The runtime attributes + * the caller itself, so OIDC tokens are omitted; only an SDK key is included + * when configured. + */ +function getRuntimeIngestHeaders( + options: IngestOptions, + flushReason: FlushReason, +): Record { + return { + 'Content-Type': 'application/json', + ...(options.auth.sdkKey + ? { Authorization: `Bearer ${options.auth.sdkKey}` } + : null), + 'User-Agent': `VercelFlagsCore/${version}`, + [FLUSH_REASON_HEADER]: flushReason, + ...((options.metricEnvironment ?? process.env.VERCEL_ENV) + ? { + 'X-Vercel-Env': + options.metricEnvironment ?? (process.env.VERCEL_ENV as string), + } + : null), + ...(isDebugMode ? { 'x-vercel-debug-ingest': '1' } : null), + }; +} + export async function sendIngestEvents( options: IngestOptions, events: UsageEvent[], @@ -81,7 +107,7 @@ export async function sendIngestEvents( const runtimeIngest = getRuntimeIngest(); if (runtimeIngest) { - const headers = await getIngestHeaders(options, flushReason); + const headers = getRuntimeIngestHeaders(options, flushReason); // Events the runtime does not accept fall through to the HTTP transport. eventsToSend = eventsToSend.filter( (event) => !runtimeIngest({ headers, body: [event] }), diff --git a/packages/vercel-flags-core/src/utils/usage-tracker.test.ts b/packages/vercel-flags-core/src/utils/usage-tracker.test.ts index ea4b7448..b972ea26 100644 --- a/packages/vercel-flags-core/src/utils/usage-tracker.test.ts +++ b/packages/vercel-flags-core/src/utils/usage-tracker.test.ts @@ -1260,11 +1260,43 @@ describe('runtime ingest transport', () => { const { headers, body } = ingestMock.mock.calls[0]![0]; expect(headers.Authorization).toBe('Bearer test-key'); expect(headers[FLUSH_REASON_HEADER]).toBe('immediate'); + expect(headers[EVALUATING_OIDC_TOKEN_HEADER]).toBeUndefined(); expect(body).toHaveLength(1); expect(body[0]!.type).toBe('FLAG_EVALUATION'); expect(fetchMock).not.toHaveBeenCalled(); expect(waitUntilMock).not.toHaveBeenCalled(); + expect(getVercelOidcTokenMock).not.toHaveBeenCalled(); + }); + + it('omits the authorization header without an SDK key', async () => { + ingestMock.mockReturnValue(true); + + const resolveToken = vi.fn(); + const tracker = new UsageTracker({ + auth: { + sdkKey: undefined, + resolveToken, + resolveBundledDefinitionsLookup: () => + Promise.resolve({ type: 'project-id' as const, projectId: 'prj_1' }), + }, + host: 'https://example.com', + fetch: fetchMock, + }); + tracker.trackEvaluation({ + flagKey: 'my-flag', + variant: 'on', + reason: ResolutionReason.RULE_MATCH, + }); + + await vi.waitFor(() => expect(ingestMock).toHaveBeenCalledTimes(1)); + + const { headers } = ingestMock.mock.calls[0]![0]; + expect(headers.Authorization).toBeUndefined(); + expect(headers[EVALUATING_OIDC_TOKEN_HEADER]).toBeUndefined(); + expect(resolveToken).not.toHaveBeenCalled(); + expect(getVercelOidcTokenMock).not.toHaveBeenCalled(); + expect(fetchMock).not.toHaveBeenCalled(); }); it('delivers each event separately', async () => { From d31aae472b0388fe9857820eb17239519c426249 Mon Sep 17 00:00:00 2001 From: Andy Bitz Date: Wed, 16 Sep 2026 18:48:13 +0200 Subject: [PATCH 3/3] Track immediate flushes so shutdown drains the HTTP fallback --- .../src/utils/usage-tracker.test.ts | 44 ++++++++++++++++++- .../src/utils/usage-tracker.ts | 21 ++++++++- 2 files changed, 62 insertions(+), 3 deletions(-) diff --git a/packages/vercel-flags-core/src/utils/usage-tracker.test.ts b/packages/vercel-flags-core/src/utils/usage-tracker.test.ts index be1ef0d8..522aec37 100644 --- a/packages/vercel-flags-core/src/utils/usage-tracker.test.ts +++ b/packages/vercel-flags-core/src/utils/usage-tracker.test.ts @@ -1252,7 +1252,7 @@ describe('runtime ingest transport', () => { delete (globalThis as Record)[FLAGS_CONTEXT_SYMBOL]; }); - it('delivers events through the runtime without fetch or waitUntil', async () => { + it('delivers events through the runtime without fetch', async () => { ingestMock.mockReturnValue(true); const tracker = createTracker(); @@ -1271,8 +1271,12 @@ describe('runtime ingest transport', () => { expect(body).toHaveLength(1); expect(body[0]!.type).toBe('FLAG_EVALUATION'); + // The flush is registered with waitUntil, but it is already settled when + // the runtime accepted every event, so it never extends the invocation. + expect(waitUntilMock).toHaveBeenCalledTimes(1); + await waitUntilMock.mock.calls[0]![0]; + expect(fetchMock).not.toHaveBeenCalled(); - expect(waitUntilMock).not.toHaveBeenCalled(); expect(getVercelOidcTokenMock).not.toHaveBeenCalled(); }); @@ -1281,6 +1285,7 @@ describe('runtime ingest transport', () => { const resolveToken = vi.fn(); const tracker = new UsageTracker({ + waitUntil, auth: { sdkKey: undefined, resolveToken, @@ -1344,6 +1349,41 @@ describe('runtime ingest transport', () => { expect(events[0]!.type).toBe('FLAG_EVALUATION'); }); + it('drains the HTTP fallback before shutdown resolves', async () => { + ingestMock.mockReturnValue(false); + + let resolveFetch!: (response: Response | PromiseLike) => void; + fetchMock.mockImplementation( + () => new Promise((res) => (resolveFetch = res)), + ); + + const tracker = createTracker(); + tracker.trackEvaluation({ + flagKey: 'my-flag', + variant: 'on', + reason: ResolutionReason.RULE_MATCH, + }); + + await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledTimes(1)); + + let shutdownResolved = false; + const shutdown = tracker.shutdown().then(() => { + shutdownResolved = true; + }); + + // Let any settled promises run: shutdown must still be waiting on the + // in-flight fallback send. + await new Promise((res) => setTimeout(res, 0)); + expect(shutdownResolved).toBe(false); + + resolveFetch(jsonResponse({ ok: true })); + await shutdown; + + const events = getBody() as SerializedEvaluationEvent[]; + expect(events).toHaveLength(1); + expect(events[0]!.type).toBe('FLAG_EVALUATION'); + }); + it('uses the scheduler when the runtime does not provide a transport', async () => { delete (globalThis as Record)[FLAGS_CONTEXT_SYMBOL]; fetchMock.mockImplementation(() => jsonResponse({ ok: true })); diff --git a/packages/vercel-flags-core/src/utils/usage-tracker.ts b/packages/vercel-flags-core/src/utils/usage-tracker.ts index af3e72d9..4b515cd8 100644 --- a/packages/vercel-flags-core/src/utils/usage-tracker.ts +++ b/packages/vercel-flags-core/src/utils/usage-tracker.ts @@ -22,6 +22,8 @@ export class UsageTracker { private options: IngestOptions; private scheduler: Scheduler; + private waitUntil: WaitUntil; + private inflightFlushes = new Set>(); private trackedRequests = new WeakSet(); @@ -30,6 +32,7 @@ export class UsageTracker { constructor(options: IngestOptions & { waitUntil: WaitUntil }) { this.options = options; + this.waitUntil = options.waitUntil; this.scheduler = new Scheduler( (reason) => this.flushEvents(reason), options.waitUntil, @@ -47,6 +50,11 @@ export class UsageTracker { // Safety net for events tracked after the drained batch reset; if the // drained flush already sent everything this returns early (maps cleared). await this.flushEvents('shutdown'); + + // Drain immediate flushes whose events fell back to the HTTP transport; + // their events left the maps before the async send started, so the flush + // above cannot cover them. + await Promise.all([...this.inflightFlushes]); } /** @@ -110,9 +118,20 @@ export class UsageTracker { */ private requestFlush(): void { if (getRuntimeIngest()) { - void this.flushEvents('immediate').catch((error) => { + // Track the flush so shutdown() can drain it: events the runtime does + // not accept fall back to the async HTTP transport, which outlives this + // synchronous call. When the runtime accepts everything the promise is + // already settled, so neither waitUntil nor shutdown() waits on it. + const flush = this.flushEvents('immediate').catch((error) => { console.error('@vercel/flags-core: Failed to flush events:', error); }); + this.inflightFlushes.add(flush); + void flush.finally(() => this.inflightFlushes.delete(flush)); + try { + this.waitUntil?.(flush); + } catch { + // waitUntil is best-effort; shutdown() still drains the flush. + } } else { this.scheduler.scheduleFlush(); }