diff --git a/.changeset/header-driven-vercel-mode.md b/.changeset/header-driven-vercel-mode.md new file mode 100644 index 00000000..06155d32 --- /dev/null +++ b/.changeset/header-driven-vercel-mode.md @@ -0,0 +1,5 @@ +--- +"@vercel/flags-core": minor +--- + +Add a header-driven `vercel` client mode that uses timestamps from the `x-vercel-flags-config-versions` or `flags-config-versions` request header instead of streaming or polling. Reuse fresh definitions, refresh in the background when the request version is up to 10 seconds newer than the cached configuration, and block for a refresh when the gap is larger. Deduplicate concurrent refreshes, allow retries after fetch failures, and discard late responses after shutdown. diff --git a/packages/vercel-flags-core/README.md b/packages/vercel-flags-core/README.md index 3cad2db8..b2de2685 100644 --- a/packages/vercel-flags-core/README.md +++ b/packages/vercel-flags-core/README.md @@ -47,6 +47,11 @@ const client = createClient(process.env.FLAGS!, { This option is sent only to the metrics ingestion endpoint. It does not select the environment used for flag evaluation. +## Configuration version headers + +The header source reads `x-vercel-flags-config-versions` or `flags-config-versions`, +with the `x-vercel-` header taking precedence when both are present. + ## OpenFeature An OpenFeature-compatible provider is available at `@vercel/flags-core/openfeature`: diff --git a/packages/vercel-flags-core/src/controller/header-source.test.ts b/packages/vercel-flags-core/src/controller/header-source.test.ts new file mode 100644 index 00000000..46b8ec94 --- /dev/null +++ b/packages/vercel-flags-core/src/controller/header-source.test.ts @@ -0,0 +1,506 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import type { BundledDefinitions, DatafileInput } from '../types'; +import { getRequestContext } from '../utils/request-context'; +import { fetchDatafile } from './fetch-datafile'; +import { HeaderSource } from './header-source'; +import { normalizeOptions } from './normalized-options'; +import { tagData } from './tagged-data'; + +vi.mock('../utils/request-context', () => ({ getRequestContext: vi.fn() })); +vi.mock('./fetch-datafile', () => ({ fetchDatafile: vi.fn() })); + +const PROJECT_ID = 'prj_test'; +const CURRENT_TIMESTAMP = 1_700_000_000_000; +const HEADER = 'x-vercel-flags-config-versions'; + +function datafile(configUpdatedAt = CURRENT_TIMESTAMP): BundledDefinitions { + return { + projectId: PROJECT_ID, + environment: 'production', + definitions: {}, + configUpdatedAt, + digest: `digest-${configUpdatedAt}`, + revision: configUpdatedAt - CURRENT_TIMESTAMP + 1, + }; +} + +function deferred() { + let resolve!: (value: T) => void; + let reject!: (reason: Error) => void; + const promise = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise; + reject = rejectPromise; + }); + return { promise, resolve, reject }; +} + +// Drain promise continuations, without sleeps, polling, or wall-clock timestamps. +function settlePromises() { + return new Promise((resolve) => setImmediate(resolve)); +} + +function setHeader(value: string | undefined) { + vi.mocked(getRequestContext).mockReturnValue({ + ctx: {}, + headers: value === undefined ? undefined : { [HEADER]: value }, + }); +} + +function setVersion(timestamp: number) { + setHeader(`flags_${PROJECT_ID}=${timestamp}`); +} + +let source: HeaderSource; +let onData: ReturnType void>>; + +beforeEach(() => { + vi.resetAllMocks(); + setVersion(CURRENT_TIMESTAMP); + vi.mocked(fetchDatafile).mockResolvedValue( + datafile(CURRENT_TIMESTAMP + 20_000), + ); + source = new HeaderSource( + normalizeOptions({ + auth: { + resolveToken: async () => 'vf_test', + resolveBundledDefinitionsLookup: async () => ({ + type: 'project-id', + projectId: PROJECT_ID, + }), + }, + // Even an accidental call through to real fetchDatafile cannot use the network. + fetch: vi + .fn() + .mockRejectedValue(new Error('Unexpected fetch')), + stream: false, + polling: false, + buildStep: false, + }), + ); + onData = vi.fn<(data: DatafileInput) => void>(); + source.on('data', onData); +}); + +afterEach(() => { + source.stop(); +}); + +describe('HeaderSource', () => { + describe('header names', () => { + it.each([ + HEADER, + 'flags-config-versions', + ])('reads %s', async (headerName) => { + vi.mocked(getRequestContext).mockReturnValue({ + ctx: {}, + headers: { [headerName]: `flags_${PROJECT_ID}=${CURRENT_TIMESTAMP}` }, + }); + const current = tagData(datafile(), 'provided'); + + expect(source.isAvailable(PROJECT_ID)).toBe(true); + await expect(source.read(current)).resolves.toEqual([current, 'HIT']); + expect(fetchDatafile).not.toHaveBeenCalled(); + }); + + it('prefers the x-vercel header when both names are present', async () => { + vi.mocked(getRequestContext).mockReturnValue({ + ctx: {}, + headers: { + [HEADER]: `flags_${PROJECT_ID}=${CURRENT_TIMESTAMP}`, + 'flags-config-versions': `flags_${PROJECT_ID}=${CURRENT_TIMESTAMP + 20_000}`, + }, + }); + const current = tagData(datafile(), 'provided'); + + await expect(source.read(current)).resolves.toEqual([current, 'HIT']); + expect(fetchDatafile).not.toHaveBeenCalled(); + }); + + it.each([ + 'x-vercel-edge-config-versions', + 'edge-config-versions', + 'x-vercel-flags-config-version', + 'flags-config-version', + ])('ignores the obsolete header %s', async (headerName) => { + vi.mocked(getRequestContext).mockReturnValue({ + ctx: {}, + headers: { + [headerName]: `flags_${PROJECT_ID}=${CURRENT_TIMESTAMP + 20_000}`, + }, + }); + + expect(source.isAvailable(PROJECT_ID)).toBe(false); + await expect( + source.read(tagData(datafile(), 'provided')), + ).resolves.toBeUndefined(); + expect(fetchDatafile).not.toHaveBeenCalled(); + }); + }); + + describe('project-specific version header', () => { + it.each([ + `flags_other=${CURRENT_TIMESTAMP + 20_000};flags_${PROJECT_ID}=${CURRENT_TIMESTAMP}`, + `flags_${PROJECT_ID}=${CURRENT_TIMESTAMP};flags_other=${CURRENT_TIMESTAMP + 20_000}`, + `flags_${PROJECT_ID}_suffix=${CURRENT_TIMESTAMP + 20_000};flags_${PROJECT_ID}=${CURRENT_TIMESTAMP}`, + ])('selects the exact project from %s', async (header) => { + setHeader(header); + const current = tagData(datafile(), 'provided'); + + expect(source.isAvailable(PROJECT_ID)).toBe(true); + await expect(source.read(current)).resolves.toEqual([current, 'HIT']); + expect(fetchDatafile).not.toHaveBeenCalled(); + }); + + it.each([ + ['absent header', undefined], + ['empty header', ''], + ['another project', `flags_other=${CURRENT_TIMESTAMP}`], + [ + 'project prefix only', + `flags_${PROJECT_ID}_suffix=${CURRENT_TIMESTAMP}`, + ], + ['missing equals sign', `flags_${PROJECT_ID}`], + ['empty value', `flags_${PROJECT_ID}=`], + ['whitespace value', `flags_${PROJECT_ID}= `], + ['non-numeric value', `flags_${PROJECT_ID}=invalid`], + [ + 'numeric prefix with garbage', + `flags_${PROJECT_ID}=${CURRENT_TIMESTAMP}ms`, + ], + ['NaN', `flags_${PROJECT_ID}=NaN`], + ['infinite value', `flags_${PROJECT_ID}=Infinity`], + ['overflowing value', `flags_${PROJECT_ID}=1e309`], + ['negative timestamp', `flags_${PROJECT_ID}=-1`], + ])('ignores %s', async (_label, header) => { + setHeader(header); + const result = await source.read(tagData(datafile(), 'provided')); + + expect.soft(source.isAvailable(PROJECT_ID)).toBe(false); + expect.soft(result).toBeUndefined(); + expect.soft(fetchDatafile).not.toHaveBeenCalled(); + expect.soft(onData).not.toHaveBeenCalled(); + }); + }); + + describe('timestamp boundaries', () => { + it.each([ + -1, 0, + ])('returns HIT without fetching for delta %i ms', async (delta) => { + setVersion(CURRENT_TIMESTAMP + delta); + const current = tagData(datafile(), 'provided'); + + const result = await source.read(current); + + expect(result).toEqual([current, 'HIT']); + expect(result?.[0]).toBe(current); + expect(fetchDatafile).not.toHaveBeenCalled(); + expect(onData).not.toHaveBeenCalled(); + }); + + it.each([ + 1, 9_999, 10_000, + ])('returns STALE immediately and emits background data for delta %i ms', async (delta) => { + setVersion(CURRENT_TIMESTAMP + delta); + const current = tagData(datafile(), 'provided'); + const fresh = datafile(CURRENT_TIMESTAMP + delta); + const pending = deferred(); + vi.mocked(fetchDatafile).mockReturnValueOnce(pending.promise); + + const result = await source.read(current); + + expect(result).toEqual([current, 'STALE']); + expect(result?.[0]).toBe(current); + expect(fetchDatafile).toHaveBeenCalledTimes(1); + expect(onData).not.toHaveBeenCalled(); + pending.resolve(fresh); + await settlePromises(); + expect(onData).toHaveBeenCalledExactlyOnceWith(fresh); + }); + + it.each([ + 10_001, 20_000, + ])('blocks and emits fetched data with MISS for delta %i ms', async (delta) => { + setVersion(CURRENT_TIMESTAMP + delta); + const fresh = datafile(CURRENT_TIMESTAMP + delta); + const pending = deferred(); + vi.mocked(fetchDatafile).mockReturnValueOnce(pending.promise); + const settled = vi.fn(); + const read = source.read(tagData(datafile(), 'provided')); + void read.then(settled); + await settlePromises(); + + expect(settled).not.toHaveBeenCalled(); + expect(onData).not.toHaveBeenCalled(); + pending.resolve(fresh); + await expect(read).resolves.toEqual([ + { ...fresh, _origin: 'fetched' }, + 'MISS', + ]); + expect(fetchDatafile).toHaveBeenCalledTimes(1); + expect(onData).toHaveBeenCalledExactlyOnceWith(fresh); + }); + + it('accepts legacy string timestamps in current data', async () => { + const current = tagData( + { ...datafile(), configUpdatedAt: String(CURRENT_TIMESTAMP) }, + 'provided', + ); + + await expect(source.read(current)).resolves.toEqual([current, 'HIT']); + expect(fetchDatafile).not.toHaveBeenCalled(); + }); + + it('returns undefined without fetching when current data has no timestamp', async () => { + setVersion(CURRENT_TIMESTAMP + 20_000); + const current = tagData( + { ...datafile(), configUpdatedAt: undefined }, + 'provided', + ); + + await expect(source.read(current)).resolves.toBeUndefined(); + expect(fetchDatafile).not.toHaveBeenCalled(); + expect(onData).not.toHaveBeenCalled(); + }); + }); + + describe('fetch deduplication and subsequent reads', () => { + it('shares one pending blocking fetch across concurrent readers', async () => { + setVersion(CURRENT_TIMESTAMP + 20_000); + const current = tagData(datafile(), 'provided'); + const pending = deferred(); + const fresh = datafile(CURRENT_TIMESTAMP + 20_000); + vi.mocked(fetchDatafile).mockReturnValueOnce(pending.promise); + + const reads = [ + source.read(current), + source.read(current), + source.read(current), + ]; + + expect(fetchDatafile).toHaveBeenCalledTimes(1); + pending.resolve(fresh); + const results = await Promise.all(reads); + for (const result of results) { + expect(result).toEqual([{ ...fresh, _origin: 'fetched' }, 'MISS']); + } + expect(onData).toHaveBeenCalledExactlyOnceWith(fresh); + }); + + it('shares background fetches even after the STALE read has settled', async () => { + setVersion(CURRENT_TIMESTAMP + 1); + const current = tagData(datafile(), 'provided'); + const pending = deferred(); + const fresh = datafile(CURRENT_TIMESTAMP + 1); + vi.mocked(fetchDatafile).mockReturnValueOnce(pending.promise); + + const reads = await Promise.all([ + source.read(current), + source.read(current), + ]); + expect(reads).toEqual([ + [current, 'STALE'], + [current, 'STALE'], + ]); + await expect(source.read(current)).resolves.toEqual([current, 'STALE']); + expect(fetchDatafile).toHaveBeenCalledTimes(1); + pending.resolve(fresh); + await settlePromises(); + expect(onData).toHaveBeenCalledExactlyOnceWith(fresh); + }); + + it('rechecks the next request header after a HIT', async () => { + const current = tagData(datafile(), 'provided'); + await expect(source.read(current)).resolves.toEqual([current, 'HIT']); + setVersion(CURRENT_TIMESTAMP + 20_000); + + const result = await source.read(current); + + expect(result?.[1]).toBe('MISS'); + expect(fetchDatafile).toHaveBeenCalledTimes(1); + }); + + it('rechecks availability when the next request has no header', async () => { + const current = tagData(datafile(), 'provided'); + await expect(source.read(current)).resolves.toEqual([current, 'HIT']); + setHeader(undefined); + + await expect(source.read(current)).resolves.toBeUndefined(); + expect(fetchDatafile).not.toHaveBeenCalled(); + }); + + it('rechecks a request after an earlier read had no header', async () => { + const current = tagData(datafile(), 'provided'); + setHeader(undefined); + await expect(source.read(current)).resolves.toBeUndefined(); + setVersion(CURRENT_TIMESTAMP); + + await expect(source.read(current)).resolves.toEqual([current, 'HIT']); + }); + + it('uses updated current data after a background fetch', async () => { + setVersion(CURRENT_TIMESTAMP + 1); + const current = tagData(datafile(), 'provided'); + const fresh = datafile(CURRENT_TIMESTAMP + 1); + vi.mocked(fetchDatafile).mockResolvedValueOnce(fresh); + await expect(source.read(current)).resolves.toEqual([current, 'STALE']); + await settlePromises(); + const updated = tagData(fresh, 'fetched'); + + const result = await source.read(updated); + + expect(result).toEqual([updated, 'HIT']); + expect(result?.[0]).toBe(updated); + expect(fetchDatafile).toHaveBeenCalledTimes(1); + }); + + it('returns HIT rather than replaying MISS after a blocking fetch', async () => { + setVersion(CURRENT_TIMESTAMP + 20_000); + const first = await source.read(tagData(datafile(), 'provided')); + expect(first?.[1]).toBe('MISS'); + const updated = tagData(datafile(CURRENT_TIMESTAMP + 20_000), 'fetched'); + + const second = await source.read(updated); + + expect(second).toEqual([updated, 'HIT']); + expect(second?.[0]).toBe(updated); + expect(fetchDatafile).toHaveBeenCalledTimes(1); + }); + + it('starts another fetch when a later request requires newer data', async () => { + setVersion(CURRENT_TIMESTAMP + 20_000); + await source.read(tagData(datafile(), 'provided')); + const updated = tagData(datafile(CURRENT_TIMESTAMP + 20_000), 'fetched'); + const newer = datafile(CURRENT_TIMESTAMP + 40_000); + setVersion(newer.configUpdatedAt); + vi.mocked(fetchDatafile).mockResolvedValueOnce(newer); + + const result = await source.read(updated); + + expect(fetchDatafile).toHaveBeenCalledTimes(2); + expect(result).toEqual([{ ...newer, _origin: 'fetched' }, 'MISS']); + expect(onData).toHaveBeenCalledTimes(2); + }); + + it('retries after a rejected blocking fetch', async () => { + setVersion(CURRENT_TIMESTAMP + 20_000); + const current = tagData(datafile(), 'provided'); + const failure = new Error('Blocking fetch failed'); + const fresh = datafile(CURRENT_TIMESTAMP + 20_000); + vi.mocked(fetchDatafile) + .mockRejectedValueOnce(failure) + .mockResolvedValueOnce(fresh); + await expect(source.read(current)).rejects.toBe(failure); + expect(onData).not.toHaveBeenCalled(); + + await expect(source.read(current)).resolves.toEqual([ + { ...fresh, _origin: 'fetched' }, + 'MISS', + ]); + expect(fetchDatafile).toHaveBeenCalledTimes(2); + expect(onData).toHaveBeenCalledExactlyOnceWith(fresh); + }); + + it('handles background rejection and retries on a later read', async () => { + setVersion(CURRENT_TIMESTAMP + 1); + const current = tagData(datafile(), 'provided'); + const pending = deferred(); + const fresh = datafile(CURRENT_TIMESTAMP + 1); + vi.mocked(fetchDatafile) + .mockReturnValueOnce(pending.promise) + .mockResolvedValueOnce(fresh); + await expect(source.read(current)).resolves.toEqual([current, 'STALE']); + + // Do not swallow rejections from HeaderSource: Vitest must report an + // unhandled rejection if the background refresh has no error handler. + pending.reject(new Error('HeaderSource background refresh failed')); + await settlePromises(); + expect(onData).not.toHaveBeenCalled(); + await expect(source.read(current)).resolves.toEqual([current, 'STALE']); + await settlePromises(); + expect(fetchDatafile).toHaveBeenCalledTimes(2); + expect(onData).toHaveBeenCalledExactlyOnceWith(fresh); + }); + }); + + describe('stop', () => { + it.each([ + 1, 20_000, + ])('keeps a restarted fetch isolated from a late aborted fetch for delta %i ms', async (delta) => { + setVersion(CURRENT_TIMESTAMP + delta); + const current = tagData(datafile(), 'provided'); + const abandoned = deferred(); + const pending = deferred(); + const fresh = datafile(CURRENT_TIMESTAMP + delta); + vi.mocked(fetchDatafile) + .mockReturnValueOnce(abandoned.promise) + .mockReturnValueOnce(pending.promise); + const abandonedOutcome = source + .read(current) + .catch((error: unknown) => error); + await settlePromises(); + + source.stop(); + const restarted = source.read(current); + abandoned.resolve(fresh); + const outcome = await abandonedOutcome; + await settlePromises(); + + expect(onData).not.toHaveBeenCalled(); + if (delta > 10_000) { + expect(outcome).toMatchObject({ name: 'AbortError' }); + } + const concurrent = source.read(current); + expect(fetchDatafile).toHaveBeenCalledTimes(2); + pending.resolve(fresh); + await Promise.all([restarted, concurrent]); + await settlePromises(); + + expect(onData).toHaveBeenCalledExactlyOnceWith(fresh); + await expect(source.read(tagData(fresh, 'fetched'))).resolves.toEqual([ + fresh, + 'HIT', + ]); + expect(fetchDatafile).toHaveBeenCalledTimes(2); + }); + + it.each([ + 1, 20_000, + ])('aborts a pending fetch for delta %i ms', async (delta) => { + setVersion(CURRENT_TIMESTAMP + delta); + const pending = deferred(); + vi.mocked(fetchDatafile).mockReturnValueOnce(pending.promise); + const read = source.read(tagData(datafile(), 'provided')); + const outcome = read.catch(() => undefined); + await settlePromises(); + const signal = vi.mocked(fetchDatafile).mock.calls[0]?.[0].signal; + + source.stop(); + // Settle even if the transport ignores cancellation, keeping tests isolated. + pending.resolve(datafile(CURRENT_TIMESTAMP + delta)); + await outcome; + await settlePromises(); + + expect(signal).toBeInstanceOf(AbortSignal); + expect(signal?.aborted).toBe(true); + }); + + it.each([ + 1, 20_000, + ])('suppresses late data emissions for delta %i ms', async (delta) => { + setVersion(CURRENT_TIMESTAMP + delta); + const pending = deferred(); + vi.mocked(fetchDatafile).mockReturnValueOnce(pending.promise); + const read = source.read(tagData(datafile(), 'provided')); + const outcome = read.catch(() => undefined); + await settlePromises(); + expect(onData).not.toHaveBeenCalled(); + + source.stop(); + pending.resolve(datafile(CURRENT_TIMESTAMP + delta)); + await outcome; + await settlePromises(); + + expect(onData).not.toHaveBeenCalled(); + }); + }); +}); diff --git a/packages/vercel-flags-core/src/controller/header-source.ts b/packages/vercel-flags-core/src/controller/header-source.ts new file mode 100644 index 00000000..83054604 --- /dev/null +++ b/packages/vercel-flags-core/src/controller/header-source.ts @@ -0,0 +1,134 @@ +import { waitUntil } from '@vercel/functions'; +import type { BundledDefinitions, DatafileInput, Metrics } from '../types'; +import { getRequestContext } from '../utils/request-context'; +import { fetchDatafile } from './fetch-datafile'; +import type { NormalizedOptions } from './normalized-options'; +import { type TaggedData, tagData } from './tagged-data'; +import { TypedEmitter } from './typed-emitter'; + +export type HeaderSourceEvents = { + data: (data: DatafileInput) => void; +}; + +/** + * Manages a lazy pulling of flag data from the flags service using the version header. + */ +export class HeaderSource extends TypedEmitter { + private options: NormalizedOptions; + private abortController: AbortController | undefined; + private promise: Promise | undefined; + + constructor(options: NormalizedOptions) { + super(); + + this.options = options; + } + + private fetchDatafile(): Promise { + // Share only the transport work, not request-specific freshness decisions. + if (this.promise) return this.promise; + + const abortController = new AbortController(); + this.abortController = abortController; + this.promise = fetchDatafile({ + ...this.options, + signal: abortController.signal, + }) + .then((data) => { + // A transport may finish after stop() even if it ignores cancellation. + abortController.signal.throwIfAborted(); + this.emit('data', data); + return data; + }) + .finally(() => { + // An older, aborted fetch must not clear a newer request's work. + if (this.abortController === abortController) { + this.promise = undefined; + this.abortController = undefined; + } + }); + + return this.promise; + } + + private getUpdatedAtHeader(projectId: string) { + const ctx = getRequestContext(); + + const header = + ctx.headers?.['x-vercel-flags-config-versions'] ?? + ctx.headers?.['flags-config-versions']; + + if (!header) { + return; + } + + const prefix = `flags_${projectId}=`; + const value = header + .split(';') + .map((part) => part.trim()) + .find((part) => part.startsWith(prefix)) + ?.slice(prefix.length); + const timestamp = Number(value); + + return Number.isFinite(timestamp) && timestamp > 0 ? timestamp : undefined; + } + + private async resolveData( + currentData: TaggedData, + ): Promise<[TaggedData, Metrics['cacheStatus']] | undefined> { + // current datafile has no timestamp, this shouldn't happen + if (!currentData.configUpdatedAt) { + return; + } + + const updatedAtHeader = this.getUpdatedAtHeader(currentData.projectId); + if (!updatedAtHeader) { + return; + } + + const currentUpdatedAt = Number(currentData.configUpdatedAt); + + // header is older than current data + if (updatedAtHeader <= currentUpdatedAt) { + return [currentData, 'HIT']; + } + + // header is within 10 seconds of current data, we can revalidate in the background + if (updatedAtHeader <= currentUpdatedAt + 10_000) { + const pending = this.fetchDatafile(); + const signal = this.abortController?.signal; + const background = pending.catch((error) => { + if (!signal?.aborted) { + console.error('@vercel/flags-core: Header refresh failed:', error); + } + }); + + waitUntil(background); + + return [currentData, 'STALE']; + } + + const data = await this.fetchDatafile(); + + return [tagData(data, 'fetched'), 'MISS']; + } + + read( + currentData: TaggedData, + ): Promise<[TaggedData, Metrics['cacheStatus']] | undefined> { + return this.resolveData(currentData); + } + + isAvailable(projectId: string): boolean { + return !!this.getUpdatedAtHeader(projectId); + } + + /** + * Abort the current header-driven fetch and discard its pending work. + */ + stop(): void { + this.abortController?.abort(); + this.abortController = undefined; + this.promise = undefined; + } +} diff --git a/packages/vercel-flags-core/src/controller/index.ts b/packages/vercel-flags-core/src/controller/index.ts index 5f55c126..101d5713 100644 --- a/packages/vercel-flags-core/src/controller/index.ts +++ b/packages/vercel-flags-core/src/controller/index.ts @@ -11,6 +11,7 @@ import type { TrackEvaluationOptions } from '../utils/usage/flags-evaluation'; import { UsageTracker } from '../utils/usage-tracker'; import { BundledSource } from './bundled-source'; import { fetchDatafile } from './fetch-datafile'; +import { HeaderSource } from './header-source'; import { type ControllerOptions, type NormalizedOptions, @@ -57,6 +58,7 @@ type State = | 'initializing:fallback' | 'streaming' | 'polling' + | 'vercel' | 'degraded' | 'build:loading' | 'build:ready' @@ -85,6 +87,11 @@ type State = * - Uses polling exclusively * - Same fallback chains as streaming mode * + * **Runtime - vercel mode** (request context has a matching x-vercel-flags-config-versions or flags-config-versions header) + * - Uses the header value to determine if the current data is fresh + * - Revalidates in the background if the header value is within 10 seconds of the current configUpdatedAt + * - Blocking fetch if header is newer than current configUpdatedAt + * * **Runtime — offline mode** (neither stream nor polling): * - Init fallback: constructor datafile → bundled → one-time fetch → throw * - Read fallback: in-memory value → constructor datafile → bundled → one-time fetch → throw @@ -108,6 +115,7 @@ export class Controller implements ControllerInterface { private streamSource: StreamSource; private pollingSource: PollingSource; private bundledSource: BundledSource; + private headerSource: HeaderSource; // Usage tracking private usageTracker: UsageTracker; @@ -136,6 +144,8 @@ export class Controller implements ControllerInterface { readBundledDefinitions, }); + this.headerSource = new HeaderSource(this.options); + // Wire source events to state machine this.wireSourceEvents(); @@ -178,6 +188,11 @@ export class Controller implements ControllerInterface { private onPollError = (error: Error) => { console.error('@vercel/flags-core: Poll failed:', error); }; + private onFetchedData = (data: DatafileInput) => { + if (this.isNewerData(data)) { + this.data = tagData(data, 'fetched'); + } + }; // --------------------------------------------------------------------------- // Source event wiring @@ -188,8 +203,11 @@ export class Controller implements ControllerInterface { this.streamSource.on('primed', this.onStreamPrimed); this.streamSource.on('connected', this.onStreamConnected); this.streamSource.on('disconnected', this.onStreamDisconnected); + this.pollingSource.on('data', this.onPollData); this.pollingSource.on('error', this.onPollError); + + this.headerSource.on('data', this.onFetchedData); } private unwireSourceEvents(): void { @@ -197,8 +215,11 @@ export class Controller implements ControllerInterface { this.streamSource.off('primed', this.onStreamPrimed); this.streamSource.off('connected', this.onStreamConnected); this.streamSource.off('disconnected', this.onStreamDisconnected); + this.pollingSource.off('data', this.onPollData); this.pollingSource.off('error', this.onPollError); + + this.headerSource.off('data', this.onFetchedData); } // --------------------------------------------------------------------------- @@ -220,6 +241,8 @@ export class Controller implements ControllerInterface { return 'streaming'; case 'polling': return 'polling'; + case 'vercel': + return 'vercel'; default: return 'offline'; } @@ -269,7 +292,9 @@ export class Controller implements ControllerInterface { // being considered initialized, so we know we have fresh data. // For no-updates (offline), return immediately since we already have usable data. if (this.data) { - if (this.options.stream.enabled) { + if (this.headerSource.isAvailable(this.data.projectId)) { + this.transition('vercel'); + } else if (this.options.stream.enabled) { this.transition('initializing:stream'); await this.tryInitializeStream(); } else if (this.options.polling.enabled) { @@ -344,6 +369,7 @@ export class Controller implements ControllerInterface { this.unwireSourceEvents(); this.streamSource.stop(); this.pollingSource.stop(); + this.headerSource.stop(); this.data = this.options.datafile ? tagData(this.options.datafile, 'provided') : undefined; @@ -442,6 +468,14 @@ export class Controller implements ControllerInterface { } if (this.data) { + if (this.headerSource.isAvailable(this.data.projectId)) { + const result = await this.headerSource.read(this.data); + + if (result) { + return result; + } + } + const cacheStatus = this.isConnected ? 'HIT' : 'STALE'; return [this.data, cacheStatus]; } diff --git a/packages/vercel-flags-core/src/types.ts b/packages/vercel-flags-core/src/types.ts index 6b67468d..04084ff1 100644 --- a/packages/vercel-flags-core/src/types.ts +++ b/packages/vercel-flags-core/src/types.ts @@ -73,7 +73,7 @@ export type Metrics = { /** Whether the stream is currently connected */ connectionState: 'connected' | 'disconnected'; /** The current operating mode of the client */ - mode: 'streaming' | 'polling' | 'build' | 'offline'; + mode: 'streaming' | 'polling' | 'build' | 'vercel' | 'offline'; /** Time in ms for the pure flag evaluation logic (only present on EvaluationResult) */ evaluationMs?: number; }; diff --git a/packages/vercel-flags-core/src/utils/usage/flags-config-read.ts b/packages/vercel-flags-core/src/utils/usage/flags-config-read.ts index c9657343..9b9930f6 100644 --- a/packages/vercel-flags-core/src/utils/usage/flags-config-read.ts +++ b/packages/vercel-flags-core/src/utils/usage/flags-config-read.ts @@ -16,7 +16,7 @@ export interface TrackReadOptions { /** Timestamp when the config was last updated */ configUpdatedAt?: number; /** The mode the SDK is operating in */ - mode?: 'poll' | 'stream' | 'build' | 'offline'; + mode?: 'poll' | 'stream' | 'build' | 'vercel' | 'offline'; /** Revision of the config */ revision?: number; } @@ -36,7 +36,7 @@ export class FlagsConfigReadEvent implements UsageEvent { duration?: number; configUpdatedAt?: number; configOrigin?: 'in-memory' | 'embedded' | 'poll' | 'stream' | 'constructor'; - mode?: 'poll' | 'stream' | 'build' | 'offline'; + mode?: TrackReadOptions['mode']; revision?: string; environment?: string; }; diff --git a/packages/vercel-flags-core/src/vercel-mode.black-box.test.ts b/packages/vercel-flags-core/src/vercel-mode.black-box.test.ts new file mode 100644 index 00000000..4bb3a5bc --- /dev/null +++ b/packages/vercel-flags-core/src/vercel-mode.black-box.test.ts @@ -0,0 +1,419 @@ +/** Public-API coverage: real client, controller, header source and fetch helper. */ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { + type BundledDefinitions, + createClient, + type FlagsClient, +} from './index.default'; +import { setRequestContext } from './test-utils'; +import { readBundledDefinitions } from './utils/read-bundled-definitions'; + +// Only the filesystem boundary is replaced; no controller/source is mocked. +vi.mock('./utils/read-bundled-definitions', () => ({ + readBundledDefinitions: vi.fn(), +})); + +const TIMESTAMP = 1_700_000_000_000; +const PROJECT_ID = 'prj_header_test'; +const HEADER = 'x-vercel-flags-config-versions'; +const SDK_KEY = 'vf_server_header_test'; + +function datafile(timestamp = TIMESTAMP, enabled = false): BundledDefinitions { + return { + definitions: { + feature: { + environments: { production: enabled ? 1 : 0 }, + variants: [false, true], + }, + }, + segments: {}, + projectId: PROJECT_ID, + environment: 'production', + configUpdatedAt: timestamp, + digest: `digest-${timestamp}`, + revision: timestamp - TIMESTAMP + 1, + }; +} + +function deferred() { + let resolve!: (value: T) => void; + let reject!: (error: Error) => void; + const promise = new Promise((res, rej) => { + resolve = res; + reject = rej; + }); + return { promise, resolve, reject }; +} + +const clients = new Set(); +const dataFetch = vi.fn(); +const transport = vi.fn(); +let cleanupContext = () => {}; + +function setVersion(timestamp?: number) { + cleanupContext(); + cleanupContext = setRequestContext( + timestamp === undefined + ? {} + : { [HEADER]: `flags_other=1;flags_${PROJECT_ID}=${timestamp}` }, + ); +} + +function client(options: Parameters[1] = {}) { + const instance = createClient(SDK_KEY, { + datafile: datafile(), + buildStep: false, + fetch: transport, + ...options, + }); + clients.add(instance); + return instance; +} + +beforeEach(() => { + vi.useFakeTimers(); + vi.setSystemTime(TIMESTAMP); + vi.stubEnv('VERCEL_ENV', 'production'); + vi.mocked(readBundledDefinitions).mockReset(); + vi.mocked(readBundledDefinitions).mockResolvedValue({ + definitions: null, + state: 'missing-file', + }); + dataFetch.mockReset(); + dataFetch.mockRejectedValue(new Error('Unexpected datafile fetch')); + transport.mockReset(); + transport.mockImplementation((input, init) => { + const url = String(input); + if (url === 'https://flags.vercel.com/v1/datafile') { + return dataFetch(input, init); + } + if (url === 'https://flags.vercel.com/v1/ingest') { + return Promise.resolve(new Response()); + } + return Promise.reject(new Error(`Unexpected request: ${url}`)); + }); + setVersion(TIMESTAMP); +}); + +afterEach(async () => { + try { + await Promise.all([...clients].map((instance) => instance.shutdown())); + } finally { + clients.clear(); + cleanupContext(); + vi.useRealTimers(); + vi.unstubAllEnvs(); + } +}); + +describe('Vercel mode (black-box)', () => { + it.each([ + HEADER, + 'flags-config-versions', + ])('initializes and refreshes using %s', async (headerName) => { + cleanupContext(); + cleanupContext = setRequestContext({ + [headerName]: `flags_${PROJECT_ID}=${TIMESTAMP}`, + }); + const instance = client(); + + const initial = await instance.evaluate('feature'); + expect(initial.value).toBe(false); + expect(initial.metrics).toMatchObject({ + mode: 'vercel', + cacheStatus: 'HIT', + }); + expect(dataFetch).not.toHaveBeenCalled(); + + cleanupContext(); + cleanupContext = setRequestContext({ + [headerName]: `flags_${PROJECT_ID}=${TIMESTAMP + 20_000}`, + }); + dataFetch.mockResolvedValueOnce( + Response.json(datafile(TIMESTAMP + 20_000, true)), + ); + + const refreshed = await instance.evaluate('feature'); + expect(refreshed.value).toBe(true); + expect(refreshed.metrics).toMatchObject({ + mode: 'vercel', + source: 'remote', + cacheStatus: 'MISS', + }); + expect(dataFetch).toHaveBeenCalledTimes(1); + }); + + it.each([ + 'provided', + 'bundled', + ] as const)('uses fresh %s definitions without opening a stream or polling', async (origin) => { + const bundled = datafile(); + vi.mocked(readBundledDefinitions).mockResolvedValue({ + definitions: bundled, + state: 'ok', + }); + const instance = client({ + datafile: origin === 'provided' ? bundled : undefined, + }); + await instance.initialize(); + + const result = await instance.evaluate('feature'); + + expect(result.value).toBe(false); + expect(result.metrics).toMatchObject({ + mode: 'vercel', + source: origin === 'provided' ? 'in-memory' : 'embedded', + cacheStatus: 'HIT', + connectionState: 'disconnected', + }); + await vi.advanceTimersByTimeAsync(60_000); + expect(dataFetch).not.toHaveBeenCalled(); + expect( + transport.mock.calls.every(([url]) => String(url).endsWith('/v1/ingest')), + ).toBe(true); + }); + + it.each([ + undefined, + 'flags_other=1700000000000', + `flags_${PROJECT_ID}=invalid`, + ])('keeps the configured offline fallback when the matching header is unavailable: %s', async (header) => { + cleanupContext(); + cleanupContext = setRequestContext(header ? { [HEADER]: header } : {}); + const instance = client({ stream: false, polling: false }); + + const result = await instance.evaluate('feature'); + + expect(result.value).toBe(false); + expect(result.metrics).toMatchObject({ + mode: 'offline', + cacheStatus: 'STALE', + }); + expect(dataFetch).not.toHaveBeenCalled(); + }); + + it('does not enable runtime header refresh during a build', async () => { + setVersion(TIMESTAMP + 20_000); + const instance = client({ buildStep: true }); + + const result = await instance.evaluate('feature'); + + expect(result.value).toBe(false); + expect(result.metrics?.mode).toBe('build'); + expect(dataFetch).not.toHaveBeenCalled(); + }); + + it.each([ + 1, 10_000, + ])('serves stale data immediately at delta %i ms, then exposes the background update', async (delta) => { + setVersion(TIMESTAMP + delta); + const pending = deferred(); + dataFetch.mockReturnValueOnce(pending.promise); + const instance = client(); + + const first = await instance.evaluate('feature'); + expect(first.value).toBe(false); + expect(first.metrics).toMatchObject({ + mode: 'vercel', + cacheStatus: 'STALE', + }); + expect(dataFetch).toHaveBeenCalledTimes(1); + expect((await instance.evaluate('feature')).value).toBe(false); + expect(dataFetch).toHaveBeenCalledTimes(1); + + pending.resolve(Response.json(datafile(TIMESTAMP + delta, true))); + await vi.advanceTimersByTimeAsync(0); + const second = await instance.evaluate('feature'); + + expect(second.value).toBe(true); + expect(second.metrics).toMatchObject({ + source: 'remote', + cacheStatus: 'HIT', + }); + expect((await instance.getDatafile()).configUpdatedAt).toBe( + TIMESTAMP + delta, + ); + expect(dataFetch).toHaveBeenCalledTimes(1); + }); + + it('blocks beyond the 10-second boundary and shares one fetch across evaluate and bulkEvaluate', async () => { + setVersion(TIMESTAMP + 10_001); + const pending = deferred(); + dataFetch.mockReturnValueOnce(pending.promise); + const instance = client(); + const settled = vi.fn(); + const single = instance.evaluate('feature').then((result) => { + settled(); + return result; + }); + const bulk = instance.bulkEvaluate([ + { key: 'feature', defaultValue: false }, + ]); + await vi.advanceTimersByTimeAsync(0); + + expect(settled).not.toHaveBeenCalled(); + expect(dataFetch).toHaveBeenCalledTimes(1); + expect(dataFetch).toHaveBeenCalledWith( + 'https://flags.vercel.com/v1/datafile', + { + headers: expect.objectContaining({ + Authorization: `Bearer ${SDK_KEY}`, + 'X-Vercel-Env': 'production', + }), + signal: expect.any(AbortSignal), + }, + ); + pending.resolve(Response.json(datafile(TIMESTAMP + 10_001, true))); + const [result, results] = await Promise.all([single, bulk]); + + for (const evaluation of [result, results.feature]) { + expect(evaluation?.value).toBe(true); + expect(evaluation?.metrics).toMatchObject({ + mode: 'vercel', + source: 'remote', + cacheStatus: 'MISS', + }); + } + expect((await instance.evaluate('feature')).metrics?.cacheStatus).toBe( + 'HIT', + ); + expect(dataFetch).toHaveBeenCalledTimes(1); + }); + + it('rechecks new request versions instead of caching the first HIT forever', async () => { + const instance = client(); + expect((await instance.evaluate('feature')).metrics?.cacheStatus).toBe( + 'HIT', + ); + + setVersion(TIMESTAMP + 20_000); + dataFetch.mockResolvedValueOnce( + Response.json(datafile(TIMESTAMP + 20_000, true)), + ); + const second = await instance.evaluate('feature'); + expect(second.value).toBe(true); + expect(second.metrics?.cacheStatus).toBe('MISS'); + + setVersion(TIMESTAMP + 40_000); + dataFetch.mockResolvedValueOnce( + Response.json(datafile(TIMESTAMP + 40_000, false)), + ); + const third = await instance.evaluate('feature'); + expect(third.value).toBe(false); + expect(third.metrics?.cacheStatus).toBe('MISS'); + expect(dataFetch).toHaveBeenCalledTimes(2); + + setVersion(); + expect((await instance.evaluate('feature')).metrics?.cacheStatus).toBe( + 'STALE', + ); + expect(dataFetch).toHaveBeenCalledTimes(2); + }); + + it('does not make a fresh request wait on another requests blocking refresh', async () => { + const instance = client(); + await instance.initialize(); + setVersion(TIMESTAMP + 20_000); + const pending = deferred(); + dataFetch.mockReturnValueOnce(pending.promise); + const blocking = instance.evaluate('feature'); + await vi.advanceTimersByTimeAsync(0); + + setVersion(TIMESTAMP); + const hitSettled = vi.fn(); + const hit = instance.evaluate('feature').then((result) => { + hitSettled(result); + return result; + }); + await vi.advanceTimersByTimeAsync(0); + // Settle the transport even when the independence assertion fails. + const completedBeforeFetch = hitSettled.mock.calls.length; + pending.resolve(Response.json(datafile(TIMESTAMP + 20_000, true))); + const [blockingResult, hitResult] = await Promise.all([blocking, hit]); + + expect(completedBeforeFetch).toBe(1); + expect(hitResult.value).toBe(false); + expect(hitResult.metrics?.cacheStatus).toBe('HIT'); + expect(blockingResult.value).toBe(true); + expect(dataFetch).toHaveBeenCalledTimes(1); + }); + + it('retries a failed blocking refresh instead of poisoning subsequent evaluations', async () => { + setVersion(TIMESTAMP + 20_000); + dataFetch.mockResolvedValueOnce( + new Response(null, { status: 503, statusText: 'Service Unavailable' }), + ); + const instance = client(); + + const failed = await instance.evaluate('feature', false); + expect(failed.value).toBe(false); + expect(failed.reason).toBe('error'); + expect(failed.errorMessage).toContain('Service Unavailable'); + + dataFetch.mockResolvedValueOnce( + Response.json(datafile(TIMESTAMP + 20_000, true)), + ); + const recovered = await instance.evaluate('feature'); + expect(recovered.value).toBe(true); + expect(recovered.metrics?.cacheStatus).toBe('MISS'); + expect(dataFetch).toHaveBeenCalledTimes(2); + }); + + it('contains background fetch errors and retries without losing cached data', async () => { + setVersion(TIMESTAMP + 1); + const pending = deferred(); + dataFetch.mockReturnValueOnce(pending.promise); + const instance = client(); + expect((await instance.evaluate('feature')).value).toBe(false); + + pending.reject(new Error('Network unavailable')); + await vi.advanceTimersByTimeAsync(0); + dataFetch.mockResolvedValueOnce( + Response.json(datafile(TIMESTAMP + 1, true)), + ); + expect((await instance.evaluate('feature')).value).toBe(false); + await vi.advanceTimersByTimeAsync(0); + + expect((await instance.evaluate('feature')).value).toBe(true); + expect(dataFetch).toHaveBeenCalledTimes(2); + }); + + it('aborts an in-flight header refresh on shutdown', async () => { + setVersion(TIMESTAMP + 1); + const pending = deferred(); + dataFetch.mockReturnValueOnce(pending.promise); + const instance = client(); + await instance.evaluate('feature'); + const signal = dataFetch.mock.calls[0]?.[1]?.signal; + + await instance.shutdown(); + clients.delete(instance); + pending.resolve(Response.json(datafile(TIMESTAMP + 1, true))); + await vi.advanceTimersByTimeAsync(0); + + expect(signal?.aborted).toBe(true); + }); + + it('reports vercel mode in config-read telemetry', async () => { + const instance = client(); + await instance.evaluate('feature'); + await instance.shutdown(); + clients.delete(instance); + + const events = transport.mock.calls + .filter(([url]) => String(url).endsWith('/v1/ingest')) + .flatMap(([, init]) => JSON.parse(String(init?.body))); + expect(events).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + type: 'FLAGS_CONFIG_READ', + payload: expect.objectContaining({ + mode: 'vercel', + configUpdatedAt: TIMESTAMP, + cacheAction: 'NONE', + }), + }), + ]), + ); + }); +});