From 559c1e7bb14765acb234bf004c8e839fadfc2922 Mon Sep 17 00:00:00 2001 From: Andy Bitz Date: Wed, 2 Sep 2026 13:10:03 +0200 Subject: [PATCH] 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 */