From 4738c09345b1c248de1efca2c262370183a8fa2c Mon Sep 17 00:00:00 2001 From: wadii Date: Tue, 22 Sep 2026 15:53:01 +0200 Subject: [PATCH 1/6] feat: experimentation support (EXP-1..5, EVT-1..5) Parse `variant`, `reason` and `metadata.experiment` from remote identity evaluations onto `Flag`, and fix `reason` parsing, which read the wrong nesting level and was therefore always undefined. Add an opt-in events pipeline behind `enableEvents`: an `EventProcessor` buffers events, deduplicates exposures per flush window, and posts them to `{eventsApiUrl}v1/events` every 10s or at 1000 events, retrying a failed batch once before dropping it. Add `getExperimentFlag`, `trackEvent`, `trackExposureEvent` and `flushEvents` on `Flagsmith`; `close()` now flushes. Local evaluation and offline mode never carry experiment metadata and never record exposures. The engine is untouched. --- index.ts | 5 + sdk/events.ts | 296 ++++++++++++++++++ sdk/index.ts | 174 ++++++++++- sdk/models.ts | 67 +++- sdk/types.ts | 40 +++ sdk/utils.ts | 17 + tests/sdk/data/identities.json | 86 +++++- tests/sdk/events.test.ts | 395 ++++++++++++++++++++++++ tests/sdk/fetchMock.ts | 6 + tests/sdk/flagsmith-experiments.test.ts | 256 +++++++++++++++ tests/sdk/models.test.ts | 116 +++++++ tests/sdk/utils.ts | 22 ++ 12 files changed, 1475 insertions(+), 5 deletions(-) create mode 100644 sdk/events.ts create mode 100644 tests/sdk/events.test.ts create mode 100644 tests/sdk/flagsmith-experiments.test.ts create mode 100644 tests/sdk/models.test.ts diff --git a/index.ts b/index.ts index cf24b7c..117f2b6 100644 --- a/index.ts +++ b/index.ts @@ -1,12 +1,17 @@ export { AnalyticsProcessor, AnalyticsProcessorOptions, + EventProcessor, + EventProcessorOptions, + FLAG_EXPOSURE_EVENT, FlagsmithAPIError, FlagsmithClientError, EnvironmentDataPollingManager, FlagsmithCache, BaseFlag, DefaultFlag, + ExperimentMetadata, + Flag, Flags, Flagsmith } from './sdk/index.js'; diff --git a/sdk/events.ts b/sdk/events.ts new file mode 100644 index 0000000..2387450 --- /dev/null +++ b/sdk/events.ts @@ -0,0 +1,296 @@ +import { pino, Logger } from 'pino'; +import { Fetch, FlagsmithTraitValue, FlagsmithValue } from './types.js'; +import { delay, getUserAgent } from './utils.js'; +import { SDK_VERSION } from './version.js'; + +/** The only `$`-prefixed event name an SDK is allowed to send. **/ +export const FLAG_EXPOSURE_EVENT = '$flag_exposure'; + +/** URL of Flagsmith's public events API. **/ +export const DEFAULT_EVENTS_API_URL = 'https://events.api.flagsmith.com/'; + +const EVENTS_ENDPOINT = 'v1/events'; + +/** Number of buffered events that triggers a flush without waiting for the timer. **/ +const DEFAULT_MAX_BUFFER = 1000; + +/** Duration in milliseconds between two automatic flushes. **/ +const DEFAULT_FLUSH_INTERVAL_MS = 10000; + +const DEFAULT_REQUEST_TIMEOUT_MS = 3000; + +/** Duration in milliseconds to wait before retrying a failed batch. **/ +const DEFAULT_RETRY_BACKOFF_MS = 1000; + +/** How many times a single batch is posted before it is dropped. **/ +const MAX_ATTEMPTS = 2; + +export interface EventProcessorOptions { + /** Client-side or server-side key of the environment that events will be recorded for. **/ + environmentKey: string; + /** {@link fetch} implementation to use for API requests. **/ + fetch: Fetch; + /** URL of the Flagsmith events API. Defaults to {@link DEFAULT_EVENTS_API_URL}. **/ + eventsApiUrl?: string; + /** Number of buffered events that triggers a flush. Defaults to {@link DEFAULT_MAX_BUFFER}. **/ + maxBuffer?: number; + /** Duration in milliseconds between automatic flushes. 0 disables the timer. Defaults to {@link DEFAULT_FLUSH_INTERVAL_MS}. **/ + flushInterval?: number; + /** Duration in milliseconds to wait for API requests to complete before timing out. Defaults to {@link DEFAULT_REQUEST_TIMEOUT_MS}. **/ + requestTimeoutMs?: number; + /** Duration in milliseconds to wait before retrying a failed batch. Defaults to {@link DEFAULT_RETRY_BACKOFF_MS}. **/ + retryBackoffMs?: number; + logger?: Logger; +} + +/** + * A single event as sent to the Flagsmith events API. + */ +export interface FlagsmithEvent { + /** The event name, e.g. `purchase` or {@link FLAG_EXPOSURE_EVENT}. **/ + event: string; + /** The feature this event relates to. Required for {@link FLAG_EXPOSURE_EVENT}. **/ + feature_name: string | null; + /** The identity this event relates to, used to reconcile exposures with conversions. **/ + identifier: string | null; + /** The event's value, always stringified by the SDK. **/ + value: string | null; + /** A flat map of resolved trait values. **/ + traits: Record | null; + /** Caller metadata, merged with the SDK version. **/ + metadata: Record; + /** Epoch milliseconds at which this event was buffered. **/ + timestamp: number; +} + +/** + * Buffers events and posts them to the Flagsmith events API. + * + * Events are flushed every {@link EventProcessorOptions.flushInterval} milliseconds, as soon as + * {@link EventProcessorOptions.maxBuffer} events are buffered, or by calling {@link flush}. + * + * Exposure events are deduplicated within a flush window. A batch that cannot be posted is retried + * once and then dropped: recording events must never fail the calling application, so all errors + * are logged and swallowed. + * @see https://docs.flagsmith.com/advanced-use/experimentation + */ +export class EventProcessor { + private eventsUrl: string; + private environmentKey: string; + private customFetch: Fetch; + private maxBuffer: number; + private flushInterval: number; + private requestTimeoutMs: number; + private retryBackoffMs: number; + private logger: Logger; + + private buffer: FlagsmithEvent[] = []; + private seenExposures: Set = new Set(); + private inFlight: Set> = new Set(); + private interval?: NodeJS.Timeout; + + constructor(opts: EventProcessorOptions) { + const eventsApiUrl = opts.eventsApiUrl || DEFAULT_EVENTS_API_URL; + this.eventsUrl = + (eventsApiUrl.endsWith('/') ? eventsApiUrl : `${eventsApiUrl}/`) + EVENTS_ENDPOINT; + this.environmentKey = opts.environmentKey; + this.customFetch = opts.fetch; + this.maxBuffer = opts.maxBuffer ?? DEFAULT_MAX_BUFFER; + this.flushInterval = opts.flushInterval ?? DEFAULT_FLUSH_INTERVAL_MS; + this.requestTimeoutMs = opts.requestTimeoutMs ?? DEFAULT_REQUEST_TIMEOUT_MS; + this.retryBackoffMs = opts.retryBackoffMs ?? DEFAULT_RETRY_BACKOFF_MS; + this.logger = opts.logger || pino(); + } + + /** + * Buffer a custom event. + */ + trackEvent(args: { + event: string; + identifier: string | null; + value: FlagsmithValue; + traits: Record | null; + metadata: Record | null; + }): void { + this.bufferEvent({ + event: args.event, + featureName: null, + identifier: args.identifier, + value: args.value, + traits: args.traits, + metadata: args.metadata + }); + } + + /** + * Buffer a {@link FLAG_EXPOSURE_EVENT} for a feature. + * + * Exposures that are identical to one already buffered in this flush window are discarded. + */ + trackExposureEvent(args: { + featureName: string; + identifier: string; + value: FlagsmithValue; + traits: Record | null; + metadata: Record | null; + }): void { + this.bufferEvent({ + event: FLAG_EXPOSURE_EVENT, + featureName: args.featureName, + identifier: args.identifier, + value: args.value, + traits: args.traits, + metadata: args.metadata + }); + } + + /** + * Post all buffered events to the Flagsmith events API. + * + * Resolves once every in-flight batch has been posted or dropped, including batches started by + * the flush timer or by reaching {@link EventProcessorOptions.maxBuffer}. + */ + async flush(): Promise { + try { + const events = this.buffer; + this.buffer = []; + this.seenExposures.clear(); + + if (events.length) { + const request = this.postEvents(events); + this.inFlight.add(request); + request.finally(() => this.inFlight.delete(request)); + } + + while (this.inFlight.size) { + await Promise.all([...this.inFlight]); + } + } catch (error) { + // Flushing events must never throw into the calling application. + this.logger.warn(error, 'Failed to flush events to the Flagsmith events API.'); + } + } + + /** + * Start flushing events every {@link EventProcessorOptions.flushInterval} milliseconds. + * + * The timer is unref'd, so it never keeps the process alive on its own. + */ + start(): void { + if (this.interval || this.flushInterval <= 0) { + return; + } + this.interval = setInterval(() => { + this.flush(); + }, this.flushInterval); + this.interval.unref?.(); + } + + /** + * Stop the flush timer and post any remaining events. + */ + async stop(): Promise { + if (this.interval) { + clearInterval(this.interval); + this.interval = undefined; + } + return this.flush(); + } + + private bufferEvent(args: { + event: string; + featureName: string | null; + identifier: string | null; + value: FlagsmithValue; + traits: Record | null; + metadata: Record | null; + }): void { + try { + const value = args.value != null ? String(args.value) : null; + const metadata: Record = { + ...(args.metadata ?? {}), + sdk_version: SDK_VERSION + }; + + if (args.event === FLAG_EXPOSURE_EVENT) { + const key = JSON.stringify([ + args.event, + args.featureName, + args.identifier, + value, + metadata['experiment_id'] ?? null + ]); + if (this.seenExposures.has(key)) { + this.logger.debug( + `Skipping duplicate exposure for feature "${args.featureName}" in this flush window.` + ); + return; + } + this.seenExposures.add(key); + } + + this.buffer.push({ + event: args.event, + feature_name: args.featureName, + identifier: args.identifier, + value: value, + traits: args.traits, + metadata: metadata, + timestamp: Date.now() + }); + + if (this.buffer.length >= this.maxBuffer) { + this.flush(); + } + } catch (error) { + // Recording an event must never throw into the calling application. + this.logger.warn(error, `Failed to record the "${args.event}" event.`); + } + } + + /** + * Post a batch, retrying once on a network error or a 5xx response before dropping it. + * + * A batch is never re-queued into the live buffer: a permanently rejected batch would otherwise + * be retried forever. + */ + private async postEvents(events: FlagsmithEvent[]): Promise { + for (let attempt = 1; attempt <= MAX_ATTEMPTS; attempt++) { + let reason = 'unknown error'; + try { + const response = await this.customFetch(this.eventsUrl, { + method: 'POST', + body: JSON.stringify({ events: events }), + signal: AbortSignal.timeout(this.requestTimeoutMs), + headers: { + 'Content-Type': 'application/json; charset=utf-8', + 'X-Environment-Key': this.environmentKey, + // The events pipeline reads the SDK language and version from this header. + 'Flagsmith-SDK-User-Agent': getUserAgent(), + 'User-Agent': getUserAgent() + } + }); + if (response.status >= 200 && response.status < 300) { + return; + } + if (response.status < 500) { + this.logger.warn( + `Flagsmith events API rejected ${events.length} events with status ${response.status}. Dropping them.` + ); + return; + } + reason = `status ${response.status}`; + } catch (error) { + reason = String(error); + } + + if (attempt === MAX_ATTEMPTS) { + this.logger.warn( + `Failed to post ${events.length} events to the Flagsmith events API (${reason}). Dropping them.` + ); + return; + } + await delay(this.retryBackoffMs); + } + } +} diff --git a/sdk/index.ts b/sdk/index.ts index 72d5f9d..0808140 100644 --- a/sdk/index.ts +++ b/sdk/index.ts @@ -6,7 +6,8 @@ import { ANALYTICS_ENDPOINT, AnalyticsProcessor } from './analytics.js'; import { BaseOfflineHandler } from './offline_handlers.js'; import { FlagsmithAPIError, FlagsmithClientError } from './errors.js'; -import { DefaultFlag, Flags } from './models.js'; +import { EventProcessor, FLAG_EXPOSURE_EVENT } from './events.js'; +import { DefaultFlag, Flag, Flags } from './models.js'; import { EnvironmentDataPollingManager } from './polling_manager.js'; import { Deferred, @@ -14,6 +15,7 @@ import { generateIdentityCacheKey, getUserAgent, isTraitConfig, + resolveTraitValues, retryFetch } from './utils.js'; import { @@ -28,6 +30,7 @@ import { FlagsmithCache, FlagsmithConfig, FlagsmithTraitValue, + FlagsmithValue, TraitConfig } from './types.js'; import { pino, Logger } from 'pino'; @@ -37,7 +40,14 @@ import { EvaluationContextWithMetadata } from '../flagsmith-engine/evaluation/mo export { AnalyticsProcessor, AnalyticsProcessorOptions } from './analytics.js'; export { FlagsmithAPIError, FlagsmithClientError } from './errors.js'; -export { BaseFlag, DefaultFlag, Flags } from './models.js'; +export { + DEFAULT_EVENTS_API_URL, + EventProcessor, + EventProcessorOptions, + FLAG_EXPOSURE_EVENT, + FlagsmithEvent +} from './events.js'; +export { BaseFlag, DefaultFlag, ExperimentMetadata, Flag, Flags } from './models.js'; export { EnvironmentDataPollingManager } from './polling_manager.js'; export { FlagsmithCache, FlagsmithConfig } from './types.js'; @@ -95,6 +105,7 @@ export class Flagsmith { private cache?: FlagsmithCache; private onEnvironmentChange: (error: Error | null, result?: EnvironmentModel) => void; private analyticsProcessor?: AnalyticsProcessor; + private eventProcessor?: EventProcessor; private logger: Logger; private customFetch: Fetch; private readonly requestRetryDelayMilliseconds: number; @@ -175,6 +186,22 @@ export class Flagsmith { logger: this.logger }); } + + if (data.eventProcessorConfig && !data.enableEvents) { + throw new Error('ValueError: eventProcessorConfig requires enableEvents: true.'); + } + + if (data.enableEvents) { + this.eventProcessor = new EventProcessor({ + ...(data.eventProcessorConfig ?? {}), + environmentKey: this.environmentKey, + fetch: this.customFetch, + requestTimeoutMs: + data.eventProcessorConfig?.requestTimeoutMs ?? this.requestTimeoutMs, + logger: this.logger + }); + this.eventProcessor.start(); + } } } /** @@ -253,6 +280,145 @@ export class Flagsmith { } } + /** + * Get a single flag for a given identity, recording a {@link FLAG_EXPOSURE_EVENT} if that + * identity is enrolled in a running experiment on the feature. + * + * Flags are fetched exactly as {@link getIdentityFlags} fetches them, so traits are upserted and + * the identity cache is used when one is configured. + * + * Experiment metadata is only produced by remote identity evaluation. Local evaluation and + * offline mode never carry it, so no exposure is ever recorded for them. + * + * @param featureName the name of the feature to evaluate. + * @param identifier a unique identifier for the identity in the current environment. + * @param traits? a dictionary of traits to add / update on the identity in Flagsmith. + * @returns the {@link Flag} for the given feature, or a {@link DefaultFlag} if it was not found. + * @throws if {@link FlagsmithConfig.enableEvents} is not set. + */ + async getExperimentFlag( + featureName: string, + identifier: string, + traits?: { [key: string]: FlagsmithTraitValue | TraitConfig } + ): Promise { + if (!this.eventProcessor) { + throw new Error('ValueError: enableEvents must be true to use getExperimentFlag.'); + } + + const flags = await this.getIdentityFlags(identifier, traits); + const flag = flags.getFlag(featureName) as Flag | DefaultFlag; + + if (!(flag instanceof Flag)) { + this.logger.debug( + `Not recording an exposure for "${featureName}": the feature was not found.` + ); + return flag; + } + if (!flag.enabled) { + this.logger.debug( + `Not recording an exposure for "${featureName}": the feature is disabled.` + ); + return flag; + } + if (!flag.experiment?.inExperiment) { + this.logger.debug( + `Not recording an exposure for "${featureName}": the identity is not enrolled in an experiment.` + ); + return flag; + } + + this.trackExposureEvent(featureName, { + identifier: identifier, + value: flag.variant, + traits: traits, + metadata: { experiment_id: flag.experiment.id } + }); + + return flag; + } + + /** + * Record a custom event, e.g. a conversion to reconcile with experiment exposures. + * + * @param event the event name. Names starting with `$` are reserved for Flagsmith. + * @param opts? the identity, value, traits and metadata to record alongside the event. + * @throws if {@link FlagsmithConfig.enableEvents} is not set, or if `event` starts with `$`. + */ + trackEvent( + event: string, + opts?: { + identifier?: string; + value?: FlagsmithValue; + traits?: { [key: string]: FlagsmithTraitValue | TraitConfig }; + metadata?: Record; + } + ): void { + if (!this.eventProcessor) { + throw new Error('ValueError: enableEvents must be true to track events.'); + } + if (event.startsWith('$')) { + throw new Error( + `ValueError: event names starting with "$" are reserved; use trackExposureEvent to record "${FLAG_EXPOSURE_EVENT}".` + ); + } + + this.eventProcessor.trackEvent({ + event: event, + identifier: opts?.identifier ?? null, + value: opts?.value ?? null, + traits: resolveTraitValues(opts?.traits), + metadata: opts?.metadata ?? null + }); + } + + /** + * Record a {@link FLAG_EXPOSURE_EVENT} for a feature. + * + * {@link getExperimentFlag} records exposures on its own; use this method to record one for a + * flag that was evaluated elsewhere. + * + * @param featureName the name of the feature the identity was exposed to. + * @param opts the identity, value, traits and metadata to record alongside the exposure. + * @throws if {@link FlagsmithConfig.enableEvents} is not set. + */ + trackExposureEvent( + featureName: string, + opts: { + identifier: string; + value?: FlagsmithValue; + traits?: { [key: string]: FlagsmithTraitValue | TraitConfig }; + metadata?: Record; + } + ): void { + if (!this.eventProcessor) { + throw new Error('ValueError: enableEvents must be true to track events.'); + } + if (!opts.identifier) { + this.logger.warn( + `Not recording an exposure for "${featureName}": an exposure requires an identifier to reconcile with conversion events.` + ); + return; + } + + this.eventProcessor.trackExposureEvent({ + featureName: featureName, + identifier: opts.identifier, + value: opts.value ?? null, + traits: resolveTraitValues(opts.traits), + metadata: opts.metadata ?? null + }); + } + + /** + * Send all buffered events to the Flagsmith events API now. + * + * Resolves once every in-flight batch has been posted or dropped, or immediately if + * {@link FlagsmithConfig.enableEvents} is not set. + */ + async flushEvents(): Promise { + await this.eventProcessor?.flush(); + } + /** * Get the segments for the current environment for a given identity. Will also upsert all traits to the Flagsmith API for future evaluations. Providing a @@ -336,8 +502,12 @@ export class Flagsmith { } } + /** + * Stop polling the environment and send any buffered events. + */ async close() { this.environmentDataPollingManager?.stop(); + await this.eventProcessor?.stop(); } private async getJSONResponse( diff --git a/sdk/models.ts b/sdk/models.ts index 1cce57e..c4bb51b 100644 --- a/sdk/models.ts +++ b/sdk/models.ts @@ -42,6 +42,57 @@ export class DefaultFlag extends BaseFlag { } } +/** + * A running experiment on a feature, as reported by a remote identity evaluation. + */ +export class ExperimentMetadata { + /** + * An identifier for this experiment, unique in a single Flagsmith installation. + */ + id: number; + /** + * The name of this experiment, as shown in the Flagsmith dashboard. + */ + name: string; + /** + * Whether this identity is enrolled in the experiment. A {@link Flag.variant} alone cannot tell: + * an identity outside the experiment's rollout is still bucketed into a variant. + */ + inExperiment: boolean; + + constructor(params: { id: number; name: string; inExperiment: boolean }) { + this.id = params.id; + this.name = params.name; + this.inExperiment = params.inExperiment; + } + + /** + * Build an {@link ExperimentMetadata} from the `metadata` object of an API flag. + * + * @param metadata The `metadata` value of an API flag. Keys other than `experiment` are ignored. + * @returns `undefined` if no usable experiment is present. + */ + static fromAPIMetadata(metadata: unknown): ExperimentMetadata | undefined { + if (!metadata || typeof metadata !== 'object') { + return undefined; + } + const experiment = (metadata as { [key: string]: any })['experiment']; + if (!experiment || typeof experiment !== 'object') { + return undefined; + } + const id = experiment['id']; + const name = experiment['name']; + if (id === undefined || id === null || name === undefined || name === null) { + return undefined; + } + return new ExperimentMetadata({ + id: id, + name: name, + inExperiment: !!experiment['in_experiment'] + }); + } +} + /** * A Flagsmith feature retrieved from a successful flag evaluation. */ @@ -58,6 +109,14 @@ export class Flag extends BaseFlag { * The reason for this feature, unique per Flagsmith project. */ reason?: string; + /** + * The variant key this identity was bucketed into. Set by remote evaluation only. + */ + variant?: string; + /** + * The experiment running on this feature, if any. Set by remote identity evaluation only. + */ + experiment?: ExperimentMetadata; constructor(params: { value: FlagValue; @@ -66,11 +125,15 @@ export class Flag extends BaseFlag { featureId: number; featureName: string; reason?: string; + variant?: string; + experiment?: ExperimentMetadata; }) { super(params.value, params.enabled, !!params.isDefault); this.featureId = params.featureId; this.featureName = params.featureName; this.reason = params.reason; + this.variant = params.variant; + this.experiment = params.experiment; } static fromFeatureStateModel( @@ -91,7 +154,9 @@ export class Flag extends BaseFlag { value: flagData['feature_state_value'] ?? flagData['value'], featureId: flagData['feature']['id'], featureName: flagData['feature']['name'], - reason: flagData['feature']['reason'] + reason: flagData['reason'], + variant: flagData['variant'], + experiment: ExperimentMetadata.fromAPIMetadata(flagData['metadata']) }); } } diff --git a/sdk/types.ts b/sdk/types.ts index ac801b3..1fda34f 100644 --- a/sdk/types.ts +++ b/sdk/types.ts @@ -117,6 +117,46 @@ export interface FlagsmithConfig { * If {@link offlineMode} is enabled, this handler is used to calculate the values of all flags. */ offlineHandler?: BaseOfflineHandler; + /** + * If enabled, the client will buffer experiment exposures and custom events and periodically + * send them to the Flagsmith events API. + * + * Required by {@link Flagsmith.getExperimentFlag}, {@link Flagsmith.trackEvent} and + * {@link Flagsmith.trackExposureEvent}. + * + * @default false + */ + enableEvents?: boolean; + /** + * Tuning options for the events pipeline. Requires {@link enableEvents}, and throws when + * provided without it. + */ + eventProcessorConfig?: { + /** + * The Flagsmith events API URL. Set this if you are not using Flagsmith's public service. + * + * @default https://events.api.flagsmith.com/ + */ + eventsApiUrl?: string; + /** + * The number of buffered events that triggers a flush without waiting for the next + * {@link flushInterval}. + * + * @default 1000 + */ + maxBuffer?: number; + /** + * The time, in milliseconds, between automatic flushes. 0 disables the timer. + * + * @default 10000 + */ + flushInterval?: number; + /** + * The events API request timeout duration, in milliseconds. Defaults to + * {@link requestTimeoutSeconds}. + */ + requestTimeoutMs?: number; + }; } /** diff --git a/sdk/utils.ts b/sdk/utils.ts index 0b838ff..6f3df02 100644 --- a/sdk/utils.ts +++ b/sdk/utils.ts @@ -42,6 +42,23 @@ export function generateIdentitiesData(identifier: string, traits: Traits, trans }; } +/** + * Unwrap any {@link TraitConfig} in a trait map to a flat map of trait values. + * + * @param traits The traits to resolve, in either supported format. + * @returns `null` if no traits were given, so that the value can be sent as-is to the events API. + */ +export function resolveTraitValues(traits?: Traits): { [key: string]: FlagsmithTraitValue } | null { + if (!traits) { + return null; + } + const resolved: { [key: string]: FlagsmithTraitValue } = {}; + for (const [key, value] of Object.entries(traits)) { + resolved[key] = isTraitConfig(value) ? value.value : value; + } + return Object.keys(resolved).length ? resolved : null; +} + export function generateIdentityCacheKey(identifier: string, traits?: Traits): string { if (!traits || Object.keys(traits).length === 0) { return `flags-${identifier}`; diff --git a/tests/sdk/data/identities.json b/tests/sdk/data/identities.json index 1d9c679..2acfc00 100644 --- a/tests/sdk/data/identities.json +++ b/tests/sdk/data/identities.json @@ -16,14 +16,96 @@ "initial_value": null, "description": null, "default_enabled": false, - "type": "STANDARD", + "type": "MULTIVARIATE", "project": 1 }, "feature_state_value": "some-value", "enabled": true, "environment": 1, "identity": null, + "feature_segment": null, + "variant": "treatment", + "reason": "SPLIT; weight=70.0", + "metadata": { + "experiment": { + "id": 167, + "name": "some_experiment", + "in_experiment": true + }, + "some_other_metadata_key": "ignored" + } + }, + { + "id": 2, + "feature": { + "id": 2, + "name": "not_enrolled_feature", + "created_date": "2019-08-27T14:53:45.698555Z", + "initial_value": null, + "description": null, + "default_enabled": false, + "type": "MULTIVARIATE", + "project": 1 + }, + "feature_state_value": "control-value", + "enabled": true, + "environment": 1, + "identity": null, + "feature_segment": null, + "variant": "control", + "reason": "SPLIT; weight=30.0", + "metadata": { + "experiment": { + "id": 168, + "name": "some_other_experiment", + "in_experiment": false + } + } + }, + { + "id": 3, + "feature": { + "id": 3, + "name": "disabled_experiment_feature", + "created_date": "2019-08-27T14:53:45.698555Z", + "initial_value": null, + "description": null, + "default_enabled": false, + "type": "MULTIVARIATE", + "project": 1 + }, + "feature_state_value": "disabled-value", + "enabled": false, + "environment": 1, + "identity": null, + "feature_segment": null, + "variant": "treatment", + "reason": "SPLIT; weight=70.0", + "metadata": { + "experiment": { + "id": 169, + "name": "disabled_experiment", + "in_experiment": true + } + } + }, + { + "id": 4, + "feature": { + "id": 4, + "name": "no_experiment_feature", + "created_date": "2019-08-27T14:53:45.698555Z", + "initial_value": null, + "description": null, + "default_enabled": false, + "type": "STANDARD", + "project": 1 + }, + "feature_state_value": "plain-value", + "enabled": true, + "environment": 1, + "identity": null, "feature_segment": null } ] -} \ No newline at end of file +} diff --git a/tests/sdk/events.test.ts b/tests/sdk/events.test.ts new file mode 100644 index 0000000..09249f0 --- /dev/null +++ b/tests/sdk/events.test.ts @@ -0,0 +1,395 @@ +import { pino } from 'pino'; +import { FLAG_EXPOSURE_EVENT } from '../../sdk/events.js'; +import { getUserAgent } from '../../sdk/utils.js'; +import { SDK_VERSION } from '../../sdk/version.js'; +import { Deferred } from '../../sdk/utils.js'; +import { fetch } from './fetchMock.js'; +import { eventProcessor, postedEvents } from './utils.js'; + +const silentLogger = pino({ level: 'silent' }); + +afterEach(() => { + vi.useRealTimers(); +}); + +test('trackEvent buffers an event with a stringified value, the SDK version and a timestamp', async () => { + const processor = eventProcessor(); + + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: 49, + traits: { plan: 'premium' }, + metadata: { currency: 'GBP' } + }); + await processor.flush(); + + const events = postedEvents(); + expect(events).toHaveLength(1); + expect(events[0]).toMatchObject({ + event: 'purchase', + feature_name: null, + identifier: 'user-123', + value: '49', + traits: { plan: 'premium' }, + metadata: { currency: 'GBP', sdk_version: SDK_VERSION } + }); + expect(Number.isInteger(events[0].timestamp)).toBe(true); +}); + +test('trackEvent keeps a missing value null', async () => { + const processor = eventProcessor(); + + processor.trackEvent({ + event: 'purchase', + identifier: null, + value: null, + traits: null, + metadata: null + }); + await processor.flush(); + + expect(postedEvents()[0]).toMatchObject({ + identifier: null, + value: null, + traits: null, + metadata: { sdk_version: SDK_VERSION } + }); +}); + +test('the SDK version wins over caller metadata', async () => { + const processor = eventProcessor(); + + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: null, + traits: null, + metadata: { sdk_version: 'not-the-sdk-version' } + }); + await processor.flush(); + + expect(postedEvents()[0].metadata.sdk_version).toBe(SDK_VERSION); +}); + +test('trackExposureEvent buffers a $flag_exposure', async () => { + const processor = eventProcessor(); + + processor.trackExposureEvent({ + featureName: 'checkout_cta', + identifier: 'user-123', + value: 'treatment', + traits: null, + metadata: { experiment_id: 167 } + }); + await processor.flush(); + + expect(postedEvents()[0]).toMatchObject({ + event: FLAG_EXPOSURE_EVENT, + feature_name: 'checkout_cta', + identifier: 'user-123', + value: 'treatment', + metadata: { experiment_id: 167, sdk_version: SDK_VERSION } + }); +}); + +test('identical exposures are deduplicated within a flush window', async () => { + const processor = eventProcessor(); + const exposure = { + featureName: 'checkout_cta', + identifier: 'user-123', + value: 'treatment', + traits: null, + metadata: { experiment_id: 167 } + }; + + processor.trackExposureEvent(exposure); + processor.trackExposureEvent(exposure); + await processor.flush(); + + expect(postedEvents()).toHaveLength(1); +}); + +test.each([ + ['identifier', { identifier: 'user-456' }], + ['value', { value: 'control' }], + ['feature', { featureName: 'other_feature' }], + ['experiment_id', { metadata: { experiment_id: 168 } }] +])('exposures differing by %s are not deduplicated', async (_name, difference) => { + const processor = eventProcessor(); + const exposure = { + featureName: 'checkout_cta', + identifier: 'user-123', + value: 'treatment', + traits: null, + metadata: { experiment_id: 167 } + }; + + processor.trackExposureEvent(exposure); + processor.trackExposureEvent({ ...exposure, ...difference }); + await processor.flush(); + + expect(postedEvents()).toHaveLength(2); +}); + +test('an identical exposure is buffered again after a flush', async () => { + const processor = eventProcessor(); + const exposure = { + featureName: 'checkout_cta', + identifier: 'user-123', + value: 'treatment', + traits: null, + metadata: { experiment_id: 167 } + }; + + processor.trackExposureEvent(exposure); + await processor.flush(); + processor.trackExposureEvent(exposure); + await processor.flush(); + + expect(postedEvents()).toHaveLength(2); +}); + +test('custom events are never deduplicated', async () => { + const processor = eventProcessor(); + const event = { + event: 'purchase', + identifier: 'user-123', + value: '49.00', + traits: null, + metadata: null + }; + + processor.trackEvent(event); + processor.trackEvent(event); + await processor.flush(); + + expect(postedEvents()).toHaveLength(2); +}); + +test('flush posts the batch to the events endpoint', async () => { + // The endpoint is built from the events API URL whether or not it has a trailing slash. + const processor = eventProcessor({ eventsApiUrl: 'http://testUrl' }); + + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: null, + traits: null, + metadata: null + }); + await processor.flush(); + + expect(fetch).toHaveBeenCalledTimes(1); + expect(fetch).toHaveBeenCalledWith( + 'http://testUrl/v1/events', + expect.objectContaining({ + method: 'POST', + headers: { + 'Content-Type': 'application/json; charset=utf-8', + 'X-Environment-Key': 'test-key', + 'Flagsmith-SDK-User-Agent': getUserAgent(), + 'User-Agent': getUserAgent() + } + }) + ); + + const body = JSON.parse(String(fetch.mock.calls[0][1]?.body)); + expect(Object.keys(body)).toEqual(['events']); + expect(body.events).toHaveLength(1); +}); + +test('flush does not post anything when nothing is buffered', async () => { + await eventProcessor().flush(); + + expect(fetch).not.toHaveBeenCalled(); +}); + +test('the buffer is flushed as soon as it reaches maxBuffer', async () => { + const processor = eventProcessor({ maxBuffer: 2 }); + + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: null, + traits: null, + metadata: null + }); + expect(fetch).not.toHaveBeenCalled(); + + processor.trackEvent({ + event: 'purchase', + identifier: 'user-456', + value: null, + traits: null, + metadata: null + }); + await processor.flush(); + + expect(fetch).toHaveBeenCalledTimes(1); + expect(postedEvents()).toHaveLength(2); +}); + +test('flush waits for a batch posted by the flush timer', async () => { + vi.useFakeTimers(); + const deferred = new Deferred(); + fetch.mockReturnValue(deferred.promise); + + const processor = eventProcessor({ flushInterval: 10000 }); + processor.start(); + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: null, + traits: null, + metadata: null + }); + + // The timer starts a flush that never settles until the response is resolved. + await vi.advanceTimersByTimeAsync(10000); + expect(fetch).toHaveBeenCalledTimes(1); + + let flushed = false; + const flush = processor.flush().then(() => { + flushed = true; + }); + await vi.advanceTimersByTimeAsync(1); + expect(flushed).toBe(false); + + deferred.resolve(new Response(null, { status: 202 })); + await flush; + expect(flushed).toBe(true); + + await processor.stop(); +}); + +test('flush waits for a batch posted by reaching maxBuffer', async () => { + vi.useFakeTimers(); + const deferred = new Deferred(); + fetch.mockReturnValue(deferred.promise); + + const processor = eventProcessor({ maxBuffer: 1 }); + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: null, + traits: null, + metadata: null + }); + expect(fetch).toHaveBeenCalledTimes(1); + + let flushed = false; + const flush = processor.flush().then(() => { + flushed = true; + }); + await vi.advanceTimersByTimeAsync(1); + expect(flushed).toBe(false); + + deferred.resolve(new Response(null, { status: 202 })); + await flush; + expect(flushed).toBe(true); +}); + +test('a batch that fails to post is retried once after the backoff, then dropped', async () => { + vi.useFakeTimers(); + fetch.mockRejectedValue(new Error('network unreachable')); + + const processor = eventProcessor({ retryBackoffMs: 1000, logger: silentLogger }); + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: null, + traits: null, + metadata: null + }); + + const flush = processor.flush(); + await vi.advanceTimersByTimeAsync(1); + expect(fetch).toHaveBeenCalledTimes(1); + + await vi.advanceTimersByTimeAsync(1000); + await flush; + expect(fetch).toHaveBeenCalledTimes(2); + + // The batch is dropped rather than re-queued into the live buffer. + await processor.flush(); + expect(fetch).toHaveBeenCalledTimes(2); +}); + +test('a batch rejected with a 5xx is retried once, then dropped', async () => { + fetch.mockResolvedValue(new Response('downstream unavailable', { status: 503 })); + + const processor = eventProcessor({ logger: silentLogger }); + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: null, + traits: null, + metadata: null + }); + await processor.flush(); + + expect(fetch).toHaveBeenCalledTimes(2); + + await processor.flush(); + expect(fetch).toHaveBeenCalledTimes(2); +}); + +test('a batch rejected with a 4xx is dropped without retrying', async () => { + fetch.mockResolvedValue(new Response('malformed batch', { status: 400 })); + + const processor = eventProcessor({ logger: silentLogger }); + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: null, + traits: null, + metadata: null + }); + await processor.flush(); + + expect(fetch).toHaveBeenCalledTimes(1); + + await processor.flush(); + expect(fetch).toHaveBeenCalledTimes(1); +}); + +test('start does not let the flush timer keep the process alive', async () => { + const setInterval = vi.spyOn(globalThis, 'setInterval'); + const processor = eventProcessor({ flushInterval: 10000 }); + + processor.start(); + // Starting twice must not leak a second timer. + processor.start(); + + expect(setInterval).toHaveBeenCalledTimes(1); + expect(setInterval.mock.results[0].value.hasRef()).toBe(false); + + await processor.stop(); +}); + +test('start does nothing when the flush interval is disabled', () => { + const setInterval = vi.spyOn(globalThis, 'setInterval'); + + eventProcessor({ flushInterval: 0 }).start(); + + expect(setInterval).not.toHaveBeenCalled(); +}); + +test('stop clears the flush timer and posts the buffered events', async () => { + vi.useFakeTimers(); + const processor = eventProcessor({ flushInterval: 10000 }); + processor.start(); + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: null, + traits: null, + metadata: null + }); + + await processor.stop(); + + expect(vi.getTimerCount()).toBe(0); + expect(postedEvents()).toHaveLength(1); +}); diff --git a/tests/sdk/fetchMock.ts b/tests/sdk/fetchMock.ts index 6dffda9..7425937 100644 --- a/tests/sdk/fetchMock.ts +++ b/tests/sdk/fetchMock.ts @@ -20,6 +20,12 @@ export function fetchImpl(url: string, options?: RequestInit) { new Response('environment-document called without a server-side key', { status: 401 }) ); } + if (url.includes('/v1/events')) { + const events = JSON.parse(String(options?.body ?? '{}'))['events'] ?? []; + return Promise.resolve( + new Response(JSON.stringify({ accepted: events.length, rejected: [] }), { status: 202 }) + ); + } if (url.includes('/flags')) { return Promise.resolve(new Response(flagsJSON, { status: 200 })); } diff --git a/tests/sdk/flagsmith-experiments.test.ts b/tests/sdk/flagsmith-experiments.test.ts new file mode 100644 index 0000000..b9b92d7 --- /dev/null +++ b/tests/sdk/flagsmith-experiments.test.ts @@ -0,0 +1,256 @@ +import { pino } from 'pino'; +import { FLAG_EXPOSURE_EVENT } from '../../sdk/events.js'; +import { DefaultFlag, Flag } from '../../sdk/models.js'; +import { FlagsmithConfig } from '../../sdk/types.js'; +import { SDK_VERSION } from '../../sdk/version.js'; +import { fetch } from './fetchMock.js'; +import { flagsmith, postedEvents } from './utils.js'; + +vi.mock('../../sdk/polling_manager'); + +const isEsmBuild = process.env.ESM_BUILD === 'true'; + +/** A client with the events pipeline enabled and its flush timer disabled. */ +function experimentsFlagsmith(params: FlagsmithConfig = {}) { + return flagsmith({ + enableEvents: true, + ...params, + eventProcessorConfig: { flushInterval: 0, ...params.eventProcessorConfig } + }); +} + +test('eventProcessorConfig without enableEvents throws at construction', () => { + expect(() => flagsmith({ eventProcessorConfig: { maxBuffer: 10 } })).toThrow( + 'ValueError: eventProcessorConfig requires enableEvents: true.' + ); +}); + +test('getExperimentFlag throws when events are disabled', async () => { + await expect(flagsmith().getExperimentFlag('some_feature', 'identifier')).rejects.toThrow( + 'ValueError: enableEvents must be true to use getExperimentFlag.' + ); +}); + +test('trackEvent throws when events are disabled', () => { + expect(() => flagsmith().trackEvent('purchase')).toThrow( + 'ValueError: enableEvents must be true to track events.' + ); +}); + +test('trackExposureEvent throws when events are disabled', () => { + expect(() => + flagsmith().trackExposureEvent('some_feature', { identifier: 'identifier' }) + ).toThrow('ValueError: enableEvents must be true to track events.'); +}); + +test('flushEvents resolves when events are disabled', async () => { + await expect(flagsmith().flushEvents()).resolves.toBeUndefined(); + expect(postedEvents()).toHaveLength(0); +}); + +test('events are also disabled in offline mode', async () => { + const flg = flagsmith({ + offlineMode: true, + offlineHandler: { getEnvironment: () => ({}) } as any, + environmentKey: undefined, + enableEvents: true + }); + + expect(() => flg.trackEvent('purchase')).toThrow( + 'ValueError: enableEvents must be true to track events.' + ); +}); + +test('trackEvent rejects event names reserved for Flagsmith', () => { + expect(() => experimentsFlagsmith().trackEvent(FLAG_EXPOSURE_EVENT)).toThrow( + `ValueError: event names starting with "$" are reserved; use trackExposureEvent to record "${FLAG_EXPOSURE_EVENT}".` + ); +}); + +test('trackEvent records a custom event', async () => { + const flg = experimentsFlagsmith(); + + flg.trackEvent('purchase', { + identifier: 'user-123', + value: 49, + metadata: { currency: 'GBP' } + }); + await flg.flushEvents(); + + expect(postedEvents()).toEqual([ + expect.objectContaining({ + event: 'purchase', + feature_name: null, + identifier: 'user-123', + value: '49', + traits: null, + metadata: { currency: 'GBP', sdk_version: SDK_VERSION } + }) + ]); +}); + +test('trackEvent unwraps TraitConfig values', async () => { + const flg = experimentsFlagsmith(); + + flg.trackEvent('purchase', { + identifier: 'user-123', + traits: { + plan: 'premium', + age: { value: 30, transient: true } + } + }); + await flg.flushEvents(); + + expect(postedEvents()[0].traits).toEqual({ plan: 'premium', age: 30 }); +}); + +test('trackExposureEvent without an identifier logs and sends nothing', async () => { + const logger = pino({ level: 'silent' }); + const warn = vi.spyOn(logger, 'warn'); + const flg = experimentsFlagsmith({ logger }); + + flg.trackExposureEvent('some_feature', { identifier: '' }); + await flg.flushEvents(); + + expect(warn).toHaveBeenCalledWith(expect.stringContaining('requires an identifier')); + expect(postedEvents()).toHaveLength(0); +}); + +test('getExperimentFlag records an exposure for an enrolled identity', async () => { + const flg = experimentsFlagsmith(); + + const flag = (await flg.getExperimentFlag('some_feature', 'user-123')) as Flag; + + expect(flag.value).toBe('some-value'); + expect(flag.variant).toBe('treatment'); + expect(flag.reason).toBe('SPLIT; weight=70.0'); + expect(flag.experiment).toEqual({ + id: 167, + name: 'some_experiment', + inExperiment: true + }); + + await flg.flushEvents(); + expect(postedEvents()).toEqual([ + expect.objectContaining({ + event: FLAG_EXPOSURE_EVENT, + feature_name: 'some_feature', + identifier: 'user-123', + value: 'treatment', + metadata: { experiment_id: 167, sdk_version: SDK_VERSION } + }) + ]); +}); + +// Skip in ESM build: instanceof fails across module boundaries +test.skipIf(isEsmBuild)('getExperimentFlag returns a Flag for an enrolled identity', async () => { + const flg = experimentsFlagsmith(); + + expect(await flg.getExperimentFlag('some_feature', 'user-123')).toBeInstanceOf(Flag); +}); + +test('getExperimentFlag sends the resolved traits with the exposure', async () => { + const flg = experimentsFlagsmith(); + + await flg.getExperimentFlag('some_feature', 'user-123', { + plan: 'premium', + age: { value: 30, transient: true } + }); + await flg.flushEvents(); + + expect(postedEvents()[0].traits).toEqual({ plan: 'premium', age: 30 }); +}); + +test.each([ + ['the identity is outside the rollout', 'not_enrolled_feature'], + ['the feature has no experiment metadata', 'no_experiment_feature'], + ['the feature is disabled', 'disabled_experiment_feature'], + ['the feature was not found', 'missing_feature'] +])('getExperimentFlag records no exposure when %s', async (_name, featureName) => { + const flg = experimentsFlagsmith(); + + await flg.getExperimentFlag(featureName, 'user-123'); + await flg.flushEvents(); + + expect(postedEvents()).toHaveLength(0); +}); + +test('getExperimentFlag still returns the flag when the identity is outside the rollout', async () => { + const flg = experimentsFlagsmith(); + + const flag = (await flg.getExperimentFlag('not_enrolled_feature', 'user-123')) as Flag; + + expect(flag.variant).toBe('control'); + expect(flag.experiment?.inExperiment).toBe(false); +}); + +test('getExperimentFlag records no exposure for a feature served by the default flag handler', async () => { + const flg = experimentsFlagsmith({ + defaultFlagHandler: () => new DefaultFlag('some-default-value', true) + }); + + const flag = await flg.getExperimentFlag('missing_feature', 'user-123'); + await flg.flushEvents(); + + expect(flag.isDefault).toBe(true); + expect(flag.value).toBe('some-default-value'); + expect(postedEvents()).toHaveLength(0); +}); + +test('getExperimentFlag records no exposure when evaluating locally', async () => { + const flg = experimentsFlagsmith({ + environmentKey: 'ser.key', + enableLocalEvaluation: true + }); + + const flag = (await flg.getExperimentFlag('some_feature', 'user-123')) as Flag; + await flg.flushEvents(); + + expect(flag.enabled).toBe(true); + expect(flag.variant).toBeUndefined(); + expect(flag.experiment).toBeUndefined(); + expect(postedEvents()).toHaveLength(0); +}); + +test('getExperimentFlag records one exposure per identity', async () => { + const flg = experimentsFlagsmith(); + + await flg.getExperimentFlag('some_feature', 'user-123'); + await flg.getExperimentFlag('some_feature', 'user-456'); + await flg.flushEvents(); + + expect(postedEvents().map(event => event.identifier)).toEqual(['user-123', 'user-456']); +}); + +test('getExperimentFlag deduplicates repeated exposures of the same identity', async () => { + const flg = experimentsFlagsmith(); + + await flg.getExperimentFlag('some_feature', 'user-123'); + await flg.getExperimentFlag('some_feature', 'user-123'); + await flg.flushEvents(); + + expect(postedEvents()).toHaveLength(1); +}); + +test('close flushes the buffered events', async () => { + const flg = experimentsFlagsmith(); + + flg.trackEvent('purchase', { identifier: 'user-123' }); + await flg.close(); + + expect(postedEvents()).toHaveLength(1); +}); + +test('the events request carries the environment key and the SDK user agent headers', async () => { + const flg = experimentsFlagsmith(); + + flg.trackEvent('purchase', { identifier: 'user-123' }); + await flg.flushEvents(); + + const [url, options] = fetch.mock.calls.find(([url]) => String(url).includes('/v1/events'))!; + expect(url).toBe('https://events.api.flagsmith.com/v1/events'); + expect((options?.headers as Record)['X-Environment-Key']).toBe( + 'sometestfakekey' + ); + expect(JSON.parse(String(options?.body))).not.toHaveProperty('environment_key'); +}); diff --git a/tests/sdk/models.test.ts b/tests/sdk/models.test.ts new file mode 100644 index 0000000..23d18f0 --- /dev/null +++ b/tests/sdk/models.test.ts @@ -0,0 +1,116 @@ +import { ExperimentMetadata, Flag, Flags } from '../../sdk/models.js'; +import { EvaluationResultWithMetadata } from '../../flagsmith-engine/evaluation/models.js'; + +const isEsmBuild = process.env.ESM_BUILD === 'true'; + +function apiFlag(overrides: { [key: string]: any } = {}) { + return { + feature: { id: 220175, name: 'checkout_cta', type: 'MULTIVARIATE' }, + enabled: true, + feature_state_value: 'buy-now', + ...overrides + }; +} + +test('fromAPIFlag sets variant, reason and experiment when present', () => { + const flag = Flag.fromAPIFlag( + apiFlag({ + variant: 'treatment', + reason: 'SPLIT; weight=70.0', + metadata: { + experiment: { id: 167, name: 'flutter_demo_exp', in_experiment: true } + } + }) + ); + + expect(flag.featureId).toBe(220175); + expect(flag.featureName).toBe('checkout_cta'); + expect(flag.value).toBe('buy-now'); + expect(flag.variant).toBe('treatment'); + expect(flag.reason).toBe('SPLIT; weight=70.0'); + expect(flag.experiment).toEqual({ + id: 167, + name: 'flutter_demo_exp', + inExperiment: true + }); +}); + +test('fromAPIFlag leaves variant, reason and experiment undefined when absent', () => { + const flag = Flag.fromAPIFlag(apiFlag()); + + expect(flag.variant).toBeUndefined(); + expect(flag.reason).toBeUndefined(); + expect(flag.experiment).toBeUndefined(); +}); + +test('fromAPIFlag keeps in_experiment false for an identity outside the rollout', () => { + const flag = Flag.fromAPIFlag( + apiFlag({ + variant: 'control', + metadata: { experiment: { id: 167, name: 'flutter_demo_exp', in_experiment: false } } + }) + ); + + expect(flag.variant).toBe('control'); + expect(flag.experiment?.inExperiment).toBe(false); +}); + +test('fromAPIFlag ignores metadata keys other than experiment', () => { + const flag = Flag.fromAPIFlag( + apiFlag({ + metadata: { + some_other_key: { id: 1 }, + experiment: { id: 167, name: 'flutter_demo_exp', in_experiment: true } + } + }) + ); + + expect(flag.experiment).toEqual({ id: 167, name: 'flutter_demo_exp', inExperiment: true }); +}); + +test('fromAPIMetadata defaults a missing in_experiment to false', () => { + const experiment = ExperimentMetadata.fromAPIMetadata({ + experiment: { id: 167, name: 'flutter_demo_exp' } + }); + + expect(experiment).toEqual({ id: 167, name: 'flutter_demo_exp', inExperiment: false }); +}); + +test.each([ + ['undefined metadata', undefined], + ['null metadata', null], + ['a non-object metadata', 'experiment'], + ['metadata without an experiment', { some_other_key: 1 }], + ['a non-object experiment', { experiment: 'flutter_demo_exp' }], + ['an experiment without an id', { experiment: { name: 'flutter_demo_exp' } }], + ['an experiment without a name', { experiment: { id: 167 } }] +])('fromAPIMetadata returns undefined for %s', (_name, metadata) => { + expect(ExperimentMetadata.fromAPIMetadata(metadata)).toBeUndefined(); +}); + +// Skip in ESM build: instanceof fails across module boundaries +test.skipIf(isEsmBuild)('fromAPIFlag returns a Flag', () => { + expect(Flag.fromAPIFlag(apiFlag())).toBeInstanceOf(Flag); +}); + +test('fromEvaluationResult leaves variant and experiment undefined', () => { + const evaluationResult = { + flags: { + some_feature: { + name: 'some_feature', + enabled: true, + value: 'some-value', + reason: 'DEFAULT', + metadata: { id: 1 } + } + }, + segments: [] + } as unknown as EvaluationResultWithMetadata; + + const flag = Flags.fromEvaluationResult(evaluationResult).getFlag('some_feature') as Flag; + + expect(flag.value).toBe('some-value'); + expect(flag.reason).toBe('DEFAULT'); + expect(flag.variant).toBeUndefined(); + expect(flag.experiment).toBeUndefined(); +}); diff --git a/tests/sdk/utils.ts b/tests/sdk/utils.ts index d7bb61a..0d4873a 100644 --- a/tests/sdk/utils.ts +++ b/tests/sdk/utils.ts @@ -1,6 +1,7 @@ import { readFileSync } from 'fs'; import { buildEnvironmentModel } from '../../flagsmith-engine/environments/util.js'; import { AnalyticsProcessor } from '../../sdk/analytics.js'; +import { EventProcessor, EventProcessorOptions } from '../../sdk/events.js'; import Flagsmith, { FlagsmithConfig } from '../../sdk/index.js'; import { Fetch, FlagsmithCache } from '../../sdk/types.js'; import { Flags } from '../../sdk/models.js'; @@ -32,6 +33,27 @@ export function analyticsProcessor() { }); } +export function eventProcessor(params: Partial = {}) { + return new EventProcessor({ + environmentKey: 'test-key', + eventsApiUrl: 'http://testUrl/', + // Tests drive flushing explicitly unless they opt in to the timer. + flushInterval: 0, + retryBackoffMs: 0, + fetch: (url, options) => fetch(url.toString(), options), + ...params + }); +} + +/** + * The events posted to the events API by the mocked fetch, flattened across all batches. + */ +export function postedEvents(): any[] { + return fetch.mock.calls + .filter(([url]) => String(url).includes('/v1/events')) + .flatMap(([, options]) => JSON.parse(String(options?.body ?? '{}'))['events'] ?? []); +} + export function apiKey(): string { return 'sometestfakekey'; } From 9c693129a3f1f121a08990af10600d008a77bed4 Mon Sep 17 00:00:00 2001 From: wadii Date: Tue, 22 Sep 2026 16:33:37 +0200 Subject: [PATCH 2/6] fix: handle a rejected event batch promise and harden the event tests The floating `finally()` on an in-flight batch had no rejection handler, so a throwing logger would surface as an unhandled rejection and break the processor's "never throws" contract. The maxBuffer test asserted only after an explicit flush, so it passed identically without any auto-flush; it now asserts before flushing. Also pin the stringification of falsy event values. --- sdk/events.ts | 4 +++- tests/sdk/events.test.ts | 24 +++++++++++++++++++++++- 2 files changed, 26 insertions(+), 2 deletions(-) diff --git a/sdk/events.ts b/sdk/events.ts index 2387450..685679d 100644 --- a/sdk/events.ts +++ b/sdk/events.ts @@ -159,7 +159,9 @@ export class EventProcessor { if (events.length) { const request = this.postEvents(events); this.inFlight.add(request); - request.finally(() => this.inFlight.delete(request)); + // Settle both ways: a rejection here would otherwise go unhandled. + const forget = () => this.inFlight.delete(request); + request.then(forget, forget); } while (this.inFlight.size) { diff --git a/tests/sdk/events.test.ts b/tests/sdk/events.test.ts index 09249f0..240db8e 100644 --- a/tests/sdk/events.test.ts +++ b/tests/sdk/events.test.ts @@ -57,6 +57,25 @@ test('trackEvent keeps a missing value null', async () => { }); }); +test.each([ + [false, 'false'], + [0, '0'], + ['', ''] +])('trackEvent stringifies the falsy value %p', async (value, expected) => { + const processor = eventProcessor(); + + processor.trackEvent({ + event: 'purchase', + identifier: 'user-123', + value: value, + traits: null, + metadata: null + }); + await processor.flush(); + + expect(postedEvents()[0].value).toBe(expected); +}); + test('the SDK version wins over caller metadata', async () => { const processor = eventProcessor(); @@ -224,10 +243,13 @@ test('the buffer is flushed as soon as it reaches maxBuffer', async () => { traits: null, metadata: null }); - await processor.flush(); + // Posted by reaching maxBuffer, before anything asks for a flush. expect(fetch).toHaveBeenCalledTimes(1); expect(postedEvents()).toHaveLength(2); + + await processor.flush(); + expect(fetch).toHaveBeenCalledTimes(1); }); test('flush waits for a batch posted by the flush timer', async () => { From 5a8486d6b4dda24d9ae415028697f97716ff3e45 Mon Sep 17 00:00:00 2001 From: wadii Date: Wed, 23 Sep 2026 09:50:19 +0200 Subject: [PATCH 3/6] fix: bound flushEvents to batches in flight when called `flush()` looped until the in-flight set was empty, so under sustained traffic that kept reaching maxBuffer it never resolved and `close()` waited for traffic to stop. It now awaits a snapshot of the batches in flight at the time of the call, which still covers every event tracked before it, via `Promise.allSettled` so one failed batch neither rejects the promise nor stops it from waiting for the others. Correct the `getExperimentFlag` return doc: a missing feature with no default flag handler yields a disabled plain object, not a `DefaultFlag`. Pin the flush semantics, the throwing-logger path, the identity-cache exposure path and the `eventProcessorConfig` wiring with tests. --- sdk/events.ts | 40 +++--- sdk/index.ts | 8 +- tests/sdk/events.test.ts | 173 ++++++++++++++---------- tests/sdk/flagsmith-experiments.test.ts | 46 ++++++- 4 files changed, 169 insertions(+), 98 deletions(-) diff --git a/sdk/events.ts b/sdk/events.ts index 685679d..ec04f28 100644 --- a/sdk/events.ts +++ b/sdk/events.ts @@ -25,6 +25,7 @@ const DEFAULT_RETRY_BACKOFF_MS = 1000; /** How many times a single batch is posted before it is dropped. **/ const MAX_ATTEMPTS = 2; +/** Options for an {@link EventProcessor}. **/ export interface EventProcessorOptions { /** Client-side or server-side key of the environment that events will be recorded for. **/ environmentKey: string; @@ -40,6 +41,7 @@ export interface EventProcessorOptions { requestTimeoutMs?: number; /** Duration in milliseconds to wait before retrying a failed batch. Defaults to {@link DEFAULT_RETRY_BACKOFF_MS}. **/ retryBackoffMs?: number; + /** Logger for dropped batches and other failures. Defaults to a new pino logger. **/ logger?: Logger; } @@ -72,7 +74,6 @@ export interface FlagsmithEvent { * Exposure events are deduplicated within a flush window. A batch that cannot be posted is retried * once and then dropped: recording events must never fail the calling application, so all errors * are logged and swallowed. - * @see https://docs.flagsmith.com/advanced-use/experimentation */ export class EventProcessor { private eventsUrl: string; @@ -147,30 +148,27 @@ export class EventProcessor { /** * Post all buffered events to the Flagsmith events API. * - * Resolves once every in-flight batch has been posted or dropped, including batches started by - * the flush timer or by reaching {@link EventProcessorOptions.maxBuffer}. + * Resolves once every batch in flight when it was called has been posted or dropped, including + * batches started by the flush timer or by reaching {@link EventProcessorOptions.maxBuffer}. + * Every event tracked before the call has therefore been sent or dropped. Batches started after + * the call are not awaited, so sustained traffic cannot keep this promise pending. */ async flush(): Promise { - try { - const events = this.buffer; - this.buffer = []; - this.seenExposures.clear(); - - if (events.length) { - const request = this.postEvents(events); - this.inFlight.add(request); - // Settle both ways: a rejection here would otherwise go unhandled. - const forget = () => this.inFlight.delete(request); - request.then(forget, forget); - } + const events = this.buffer; + this.buffer = []; + this.seenExposures.clear(); - while (this.inFlight.size) { - await Promise.all([...this.inFlight]); - } - } catch (error) { - // Flushing events must never throw into the calling application. - this.logger.warn(error, 'Failed to flush events to the Flagsmith events API.'); + if (events.length) { + const batch = this.postEvents(events); + this.inFlight.add(batch); + // Settle both ways: a rejection here would otherwise go unhandled. + const forget = () => this.inFlight.delete(batch); + batch.then(forget, forget); } + + // allSettled, not all: one failed batch must neither reject this promise nor stop it from + // waiting for the others. + await Promise.allSettled([...this.inFlight]); } /** diff --git a/sdk/index.ts b/sdk/index.ts index 0808140..e127502 100644 --- a/sdk/index.ts +++ b/sdk/index.ts @@ -293,7 +293,9 @@ export class Flagsmith { * @param featureName the name of the feature to evaluate. * @param identifier a unique identifier for the identity in the current environment. * @param traits? a dictionary of traits to add / update on the identity in Flagsmith. - * @returns the {@link Flag} for the given feature, or a {@link DefaultFlag} if it was not found. + * @returns the {@link Flag} for the given feature. If it was not found, the result of + * {@link FlagsmithConfig.defaultFlagHandler}, or a disabled flag with `isDefault` set if there is + * no handler. * @throws if {@link FlagsmithConfig.enableEvents} is not set. */ async getExperimentFlag( @@ -412,8 +414,8 @@ export class Flagsmith { /** * Send all buffered events to the Flagsmith events API now. * - * Resolves once every in-flight batch has been posted or dropped, or immediately if - * {@link FlagsmithConfig.enableEvents} is not set. + * Resolves once every event tracked before the call has been posted or dropped, or immediately + * if {@link FlagsmithConfig.enableEvents} is not set. */ async flushEvents(): Promise { await this.eventProcessor?.flush(); diff --git a/tests/sdk/events.test.ts b/tests/sdk/events.test.ts index 240db8e..229133c 100644 --- a/tests/sdk/events.test.ts +++ b/tests/sdk/events.test.ts @@ -1,13 +1,28 @@ import { pino } from 'pino'; -import { FLAG_EXPOSURE_EVENT } from '../../sdk/events.js'; -import { getUserAgent } from '../../sdk/utils.js'; +import { EventProcessor, FLAG_EXPOSURE_EVENT } from '../../sdk/events.js'; +import { Deferred, getUserAgent } from '../../sdk/utils.js'; import { SDK_VERSION } from '../../sdk/version.js'; -import { Deferred } from '../../sdk/utils.js'; import { fetch } from './fetchMock.js'; import { eventProcessor, postedEvents } from './utils.js'; const silentLogger = pino({ level: 'silent' }); +/** Buffer a custom event that only carries an identifier. */ +function trackPurchase(processor: EventProcessor, identifier: string = 'user-123') { + processor.trackEvent({ + event: 'purchase', + identifier: identifier, + value: null, + traits: null, + metadata: null + }); +} + +/** Let every pending promise callback run. */ +function settle() { + return new Promise(resolve => setImmediate(resolve)); +} + afterEach(() => { vi.useRealTimers(); }); @@ -190,13 +205,7 @@ test('flush posts the batch to the events endpoint', async () => { // The endpoint is built from the events API URL whether or not it has a trailing slash. const processor = eventProcessor({ eventsApiUrl: 'http://testUrl' }); - processor.trackEvent({ - event: 'purchase', - identifier: 'user-123', - value: null, - traits: null, - metadata: null - }); + trackPurchase(processor); await processor.flush(); expect(fetch).toHaveBeenCalledTimes(1); @@ -227,22 +236,10 @@ test('flush does not post anything when nothing is buffered', async () => { test('the buffer is flushed as soon as it reaches maxBuffer', async () => { const processor = eventProcessor({ maxBuffer: 2 }); - processor.trackEvent({ - event: 'purchase', - identifier: 'user-123', - value: null, - traits: null, - metadata: null - }); + trackPurchase(processor); expect(fetch).not.toHaveBeenCalled(); - processor.trackEvent({ - event: 'purchase', - identifier: 'user-456', - value: null, - traits: null, - metadata: null - }); + trackPurchase(processor, 'user-456'); // Posted by reaching maxBuffer, before anything asks for a flush. expect(fetch).toHaveBeenCalledTimes(1); @@ -259,13 +256,7 @@ test('flush waits for a batch posted by the flush timer', async () => { const processor = eventProcessor({ flushInterval: 10000 }); processor.start(); - processor.trackEvent({ - event: 'purchase', - identifier: 'user-123', - value: null, - traits: null, - metadata: null - }); + trackPurchase(processor); // The timer starts a flush that never settles until the response is resolved. await vi.advanceTimersByTimeAsync(10000); @@ -291,13 +282,7 @@ test('flush waits for a batch posted by reaching maxBuffer', async () => { fetch.mockReturnValue(deferred.promise); const processor = eventProcessor({ maxBuffer: 1 }); - processor.trackEvent({ - event: 'purchase', - identifier: 'user-123', - value: null, - traits: null, - metadata: null - }); + trackPurchase(processor); expect(fetch).toHaveBeenCalledTimes(1); let flushed = false; @@ -312,24 +297,94 @@ test('flush waits for a batch posted by reaching maxBuffer', async () => { expect(flushed).toBe(true); }); +test('flush does not wait for a batch started after it was called', async () => { + const first = new Deferred(); + const second = new Deferred(); + fetch.mockReturnValueOnce(first.promise).mockReturnValueOnce(second.promise); + const processor = eventProcessor(); + + trackPurchase(processor); + let firstFlushed = false; + const firstFlush = processor.flush().then(() => { + firstFlushed = true; + }); + + // Traffic keeps arriving while the first batch is on the wire. + trackPurchase(processor, 'user-456'); + let secondFlushed = false; + const secondFlush = processor.flush().then(() => { + secondFlushed = true; + }); + + first.resolve(new Response(null, { status: 202 })); + await firstFlush; + await settle(); + expect(firstFlushed).toBe(true); + expect(secondFlushed).toBe(false); + + second.resolve(new Response(null, { status: 202 })); + await secondFlush; + expect(secondFlushed).toBe(true); +}); + +test('flush waits for every in-flight batch even when one of them rejects', async () => { + const pending = new Deferred(); + fetch + .mockRejectedValueOnce(new Error('network unreachable')) + .mockReturnValueOnce(pending.promise) + .mockRejectedValueOnce(new Error('network unreachable')); + const logger = pino({ level: 'silent' }); + // A throwing logger makes the dropped batch reject instead of resolving. + vi.spyOn(logger, 'warn').mockImplementation(() => { + throw new Error('logger failure'); + }); + const processor = eventProcessor({ maxBuffer: 1, logger }); + + // Each event reaches maxBuffer and starts its own batch. + trackPurchase(processor); + trackPurchase(processor, 'user-456'); + let flushed = false; + const flush = processor.flush().finally(() => { + flushed = true; + }); + + // The first batch has been retried and dropped, the second one is still on the wire. + await vi.waitFor(() => expect(fetch).toHaveBeenCalledTimes(3)); + await settle(); + expect(flushed).toBe(false); + + pending.resolve(new Response(null, { status: 202 })); + await flush; + expect(flushed).toBe(true); +}); + +test('a throwing logger cannot make flush reject', async () => { + fetch.mockRejectedValue(new Error('network unreachable')); + const logger = pino({ level: 'silent' }); + vi.spyOn(logger, 'warn').mockImplementation(() => { + throw new Error('logger failure'); + }); + const processor = eventProcessor({ logger }); + + trackPurchase(processor); + + // An unhandled rejection from the dropped batch would also fail this test run. + await expect(processor.flush()).resolves.toBeUndefined(); + expect(logger.warn).toHaveBeenCalled(); +}); + test('a batch that fails to post is retried once after the backoff, then dropped', async () => { vi.useFakeTimers(); fetch.mockRejectedValue(new Error('network unreachable')); const processor = eventProcessor({ retryBackoffMs: 1000, logger: silentLogger }); - processor.trackEvent({ - event: 'purchase', - identifier: 'user-123', - value: null, - traits: null, - metadata: null - }); + trackPurchase(processor); const flush = processor.flush(); - await vi.advanceTimersByTimeAsync(1); + await vi.advanceTimersByTimeAsync(999); expect(fetch).toHaveBeenCalledTimes(1); - await vi.advanceTimersByTimeAsync(1000); + await vi.advanceTimersByTimeAsync(1); await flush; expect(fetch).toHaveBeenCalledTimes(2); @@ -342,13 +397,7 @@ test('a batch rejected with a 5xx is retried once, then dropped', async () => { fetch.mockResolvedValue(new Response('downstream unavailable', { status: 503 })); const processor = eventProcessor({ logger: silentLogger }); - processor.trackEvent({ - event: 'purchase', - identifier: 'user-123', - value: null, - traits: null, - metadata: null - }); + trackPurchase(processor); await processor.flush(); expect(fetch).toHaveBeenCalledTimes(2); @@ -361,13 +410,7 @@ test('a batch rejected with a 4xx is dropped without retrying', async () => { fetch.mockResolvedValue(new Response('malformed batch', { status: 400 })); const processor = eventProcessor({ logger: silentLogger }); - processor.trackEvent({ - event: 'purchase', - identifier: 'user-123', - value: null, - traits: null, - metadata: null - }); + trackPurchase(processor); await processor.flush(); expect(fetch).toHaveBeenCalledTimes(1); @@ -402,13 +445,7 @@ test('stop clears the flush timer and posts the buffered events', async () => { vi.useFakeTimers(); const processor = eventProcessor({ flushInterval: 10000 }); processor.start(); - processor.trackEvent({ - event: 'purchase', - identifier: 'user-123', - value: null, - traits: null, - metadata: null - }); + trackPurchase(processor); await processor.stop(); diff --git a/tests/sdk/flagsmith-experiments.test.ts b/tests/sdk/flagsmith-experiments.test.ts index b9b92d7..97b3d98 100644 --- a/tests/sdk/flagsmith-experiments.test.ts +++ b/tests/sdk/flagsmith-experiments.test.ts @@ -4,7 +4,7 @@ import { DefaultFlag, Flag } from '../../sdk/models.js'; import { FlagsmithConfig } from '../../sdk/types.js'; import { SDK_VERSION } from '../../sdk/version.js'; import { fetch } from './fetchMock.js'; -import { flagsmith, postedEvents } from './utils.js'; +import { flagsmith, postedEvents, TestCache } from './utils.js'; vi.mock('../../sdk/polling_manager'); @@ -61,11 +61,14 @@ test('events are also disabled in offline mode', async () => { ); }); -test('trackEvent rejects event names reserved for Flagsmith', () => { - expect(() => experimentsFlagsmith().trackEvent(FLAG_EXPOSURE_EVENT)).toThrow( - `ValueError: event names starting with "$" are reserved; use trackExposureEvent to record "${FLAG_EXPOSURE_EVENT}".` - ); -}); +test.each([FLAG_EXPOSURE_EVENT, '$purchase'])( + 'trackEvent rejects the reserved event name %s', + event => { + expect(() => experimentsFlagsmith().trackEvent(event)).toThrow( + `ValueError: event names starting with "$" are reserved; use trackExposureEvent to record "${FLAG_EXPOSURE_EVENT}".` + ); + } +); test('trackEvent records a custom event', async () => { const flg = experimentsFlagsmith(); @@ -232,6 +235,21 @@ test('getExperimentFlag deduplicates repeated exposures of the same identity', a expect(postedEvents()).toHaveLength(1); }); +test('getExperimentFlag records an exposure for flags served from the identity cache', async () => { + const flg = experimentsFlagsmith({ cache: new TestCache() }); + + await flg.getExperimentFlag('some_feature', 'user-123'); + await flg.flushEvents(); + await flg.getExperimentFlag('some_feature', 'user-123'); + await flg.flushEvents(); + + const identityRequests = fetch.mock.calls.filter(([url]) => + String(url).includes('/identities') + ); + expect(identityRequests).toHaveLength(1); + expect(postedEvents().map(event => event.identifier)).toEqual(['user-123', 'user-123']); +}); + test('close flushes the buffered events', async () => { const flg = experimentsFlagsmith(); @@ -254,3 +272,19 @@ test('the events request carries the environment key and the SDK user agent head ); expect(JSON.parse(String(options?.body))).not.toHaveProperty('environment_key'); }); + +test('eventProcessorConfig is passed to the event processor', async () => { + const flg = flagsmith({ + enableEvents: true, + eventProcessorConfig: { eventsApiUrl: 'https://events.example.com', maxBuffer: 1 } + }); + + // Reaching maxBuffer posts the event without a flush. + flg.trackEvent('purchase', { identifier: 'user-123' }); + + expect(fetch).toHaveBeenCalledWith( + 'https://events.example.com/v1/events', + expect.objectContaining({ method: 'POST' }) + ); + await flg.close(); +}); From 87b3e16c2e90283a2797018ff811672f31de9ffe Mon Sep 17 00:00:00 2001 From: wadii Date: Wed, 23 Sep 2026 16:15:59 +0200 Subject: [PATCH 4/6] fix: send events through the client's agent and custom headers Flag and identity requests go through the configured `agent` and `customHeaders`, but the event processor received only `fetch`, so a client behind a proxy would evaluate flags and silently drop every event. Custom headers are applied first so the SDK's own headers, which the events pipeline parses for language and version, cannot be overridden. --- sdk/events.ts | 18 ++++++++++++-- sdk/index.ts | 2 ++ tests/sdk/events.test.ts | 31 +++++++++++++++++++++++++ tests/sdk/flagsmith-experiments.test.ts | 19 +++++++++++++++ 4 files changed, 68 insertions(+), 2 deletions(-) diff --git a/sdk/events.ts b/sdk/events.ts index ec04f28..925be08 100644 --- a/sdk/events.ts +++ b/sdk/events.ts @@ -1,4 +1,5 @@ import { pino, Logger } from 'pino'; +import { Dispatcher } from 'undici-types'; import { Fetch, FlagsmithTraitValue, FlagsmithValue } from './types.js'; import { delay, getUserAgent } from './utils.js'; import { SDK_VERSION } from './version.js'; @@ -31,6 +32,10 @@ export interface EventProcessorOptions { environmentKey: string; /** {@link fetch} implementation to use for API requests. **/ fetch: Fetch; + /** Custom {@link Dispatcher} to use when making HTTP requests. **/ + agent?: Dispatcher; + /** Custom headers to send with every request. The SDK's own headers take precedence. **/ + customHeaders?: { [key: string]: string }; /** URL of the Flagsmith events API. Defaults to {@link DEFAULT_EVENTS_API_URL}. **/ eventsApiUrl?: string; /** Number of buffered events that triggers a flush. Defaults to {@link DEFAULT_MAX_BUFFER}. **/ @@ -79,6 +84,8 @@ export class EventProcessor { private eventsUrl: string; private environmentKey: string; private customFetch: Fetch; + private agent?: Dispatcher; + private customHeaders?: { [key: string]: string }; private maxBuffer: number; private flushInterval: number; private requestTimeoutMs: number; @@ -96,6 +103,8 @@ export class EventProcessor { (eventsApiUrl.endsWith('/') ? eventsApiUrl : `${eventsApiUrl}/`) + EVENTS_ENDPOINT; this.environmentKey = opts.environmentKey; this.customFetch = opts.fetch; + this.agent = opts.agent; + this.customHeaders = opts.customHeaders; this.maxBuffer = opts.maxBuffer ?? DEFAULT_MAX_BUFFER; this.flushInterval = opts.flushInterval ?? DEFAULT_FLUSH_INTERVAL_MS; this.requestTimeoutMs = opts.requestTimeoutMs ?? DEFAULT_REQUEST_TIMEOUT_MS; @@ -258,18 +267,23 @@ export class EventProcessor { for (let attempt = 1; attempt <= MAX_ATTEMPTS; attempt++) { let reason = 'unknown error'; try { - const response = await this.customFetch(this.eventsUrl, { + // built-in RequestInit type doesn't have dispatcher/agent + const init: RequestInit & { dispatcher?: Dispatcher } = { + dispatcher: this.agent, method: 'POST', body: JSON.stringify({ events: events }), signal: AbortSignal.timeout(this.requestTimeoutMs), headers: { + // Custom headers first: the SDK's own headers must not be overridden. + ...(this.customHeaders ?? {}), 'Content-Type': 'application/json; charset=utf-8', 'X-Environment-Key': this.environmentKey, // The events pipeline reads the SDK language and version from this header. 'Flagsmith-SDK-User-Agent': getUserAgent(), 'User-Agent': getUserAgent() } - }); + }; + const response = await this.customFetch(this.eventsUrl, init); if (response.status >= 200 && response.status < 300) { return; } diff --git a/sdk/index.ts b/sdk/index.ts index e127502..e262278 100644 --- a/sdk/index.ts +++ b/sdk/index.ts @@ -196,6 +196,8 @@ export class Flagsmith { ...(data.eventProcessorConfig ?? {}), environmentKey: this.environmentKey, fetch: this.customFetch, + agent: this.agent, + customHeaders: this.customHeaders, requestTimeoutMs: data.eventProcessorConfig?.requestTimeoutMs ?? this.requestTimeoutMs, logger: this.logger diff --git a/tests/sdk/events.test.ts b/tests/sdk/events.test.ts index 229133c..26c798d 100644 --- a/tests/sdk/events.test.ts +++ b/tests/sdk/events.test.ts @@ -227,6 +227,37 @@ test('flush posts the batch to the events endpoint', async () => { expect(body.events).toHaveLength(1); }); +test('flush posts through the configured dispatcher', async () => { + const agent = { name: 'test-dispatcher' } as any; + const processor = eventProcessor({ agent }); + + trackPurchase(processor); + await processor.flush(); + + expect(fetch.mock.calls[0][1]).toMatchObject({ dispatcher: agent }); +}); + +test('flush sends custom headers without letting them override the SDK headers', async () => { + const processor = eventProcessor({ + customHeaders: { + 'X-Proxy-Token': 'secret', + 'Flagsmith-SDK-User-Agent': 'not-the-sdk', + 'X-Environment-Key': 'not-the-environment' + } + }); + + trackPurchase(processor); + await processor.flush(); + + expect(fetch.mock.calls[0][1]?.headers).toEqual({ + 'X-Proxy-Token': 'secret', + 'Content-Type': 'application/json; charset=utf-8', + 'X-Environment-Key': 'test-key', + 'Flagsmith-SDK-User-Agent': getUserAgent(), + 'User-Agent': getUserAgent() + }); +}); + test('flush does not post anything when nothing is buffered', async () => { await eventProcessor().flush(); diff --git a/tests/sdk/flagsmith-experiments.test.ts b/tests/sdk/flagsmith-experiments.test.ts index 97b3d98..5af75f4 100644 --- a/tests/sdk/flagsmith-experiments.test.ts +++ b/tests/sdk/flagsmith-experiments.test.ts @@ -2,6 +2,7 @@ import { pino } from 'pino'; import { FLAG_EXPOSURE_EVENT } from '../../sdk/events.js'; import { DefaultFlag, Flag } from '../../sdk/models.js'; import { FlagsmithConfig } from '../../sdk/types.js'; +import { getUserAgent } from '../../sdk/utils.js'; import { SDK_VERSION } from '../../sdk/version.js'; import { fetch } from './fetchMock.js'; import { flagsmith, postedEvents, TestCache } from './utils.js'; @@ -288,3 +289,21 @@ test('eventProcessorConfig is passed to the event processor', async () => { ); await flg.close(); }); + +test('the events request inherits the client agent and custom headers', async () => { + const agent = { name: 'test-dispatcher' } as any; + const flg = experimentsFlagsmith({ agent, customHeaders: { 'X-Proxy-Token': 'secret' } }); + + flg.trackEvent('purchase', { identifier: 'user-123' }); + await flg.flushEvents(); + + const [, options] = fetch.mock.calls.find(([url]) => String(url).includes('/v1/events'))!; + expect(options).toMatchObject({ + dispatcher: agent, + headers: expect.objectContaining({ + 'X-Proxy-Token': 'secret', + 'X-Environment-Key': 'sometestfakekey', + 'Flagsmith-SDK-User-Agent': getUserAgent() + }) + }); +}); From 56d0d174f01b6de3167041d5f5fce0423898e3db Mon Sep 17 00:00:00 2001 From: wadii Date: Wed, 23 Sep 2026 16:26:17 +0200 Subject: [PATCH 5/6] fix: drop case variants of the SDK's own event headers Header names are case-insensitive, so a custom `x-environment-key` would be combined with the SDK's `X-Environment-Key` rather than overridden by it. Reserved names are now filtered out of custom headers regardless of case. --- sdk/events.ts | 18 +++++++++++++++--- tests/sdk/events.test.ts | 4 +++- 2 files changed, 18 insertions(+), 4 deletions(-) diff --git a/sdk/events.ts b/sdk/events.ts index 925be08..784610f 100644 --- a/sdk/events.ts +++ b/sdk/events.ts @@ -26,6 +26,14 @@ const DEFAULT_RETRY_BACKOFF_MS = 1000; /** How many times a single batch is posted before it is dropped. **/ const MAX_ATTEMPTS = 2; +/** Headers the SDK sets itself. Header names are case-insensitive, so custom variants are dropped. **/ +const RESERVED_HEADERS = [ + 'content-type', + 'x-environment-key', + 'flagsmith-sdk-user-agent', + 'user-agent' +]; + /** Options for an {@link EventProcessor}. **/ export interface EventProcessorOptions { /** Client-side or server-side key of the environment that events will be recorded for. **/ @@ -85,7 +93,7 @@ export class EventProcessor { private environmentKey: string; private customFetch: Fetch; private agent?: Dispatcher; - private customHeaders?: { [key: string]: string }; + private customHeaders: { [key: string]: string }; private maxBuffer: number; private flushInterval: number; private requestTimeoutMs: number; @@ -104,7 +112,11 @@ export class EventProcessor { this.environmentKey = opts.environmentKey; this.customFetch = opts.fetch; this.agent = opts.agent; - this.customHeaders = opts.customHeaders; + this.customHeaders = Object.fromEntries( + Object.entries(opts.customHeaders ?? {}).filter( + ([name]) => !RESERVED_HEADERS.includes(name.toLowerCase()) + ) + ); this.maxBuffer = opts.maxBuffer ?? DEFAULT_MAX_BUFFER; this.flushInterval = opts.flushInterval ?? DEFAULT_FLUSH_INTERVAL_MS; this.requestTimeoutMs = opts.requestTimeoutMs ?? DEFAULT_REQUEST_TIMEOUT_MS; @@ -275,7 +287,7 @@ export class EventProcessor { signal: AbortSignal.timeout(this.requestTimeoutMs), headers: { // Custom headers first: the SDK's own headers must not be overridden. - ...(this.customHeaders ?? {}), + ...this.customHeaders, 'Content-Type': 'application/json; charset=utf-8', 'X-Environment-Key': this.environmentKey, // The events pipeline reads the SDK language and version from this header. diff --git a/tests/sdk/events.test.ts b/tests/sdk/events.test.ts index 26c798d..5cc950f 100644 --- a/tests/sdk/events.test.ts +++ b/tests/sdk/events.test.ts @@ -242,7 +242,9 @@ test('flush sends custom headers without letting them override the SDK headers', customHeaders: { 'X-Proxy-Token': 'secret', 'Flagsmith-SDK-User-Agent': 'not-the-sdk', - 'X-Environment-Key': 'not-the-environment' + // Header names are case-insensitive: a variant would be merged, not overridden. + 'x-environment-key': 'not-the-environment', + 'USER-AGENT': 'not-the-sdk' } }); From 5970dc00e06b7dd658893dad8a8cd2836f27e7ec Mon Sep 17 00:00:00 2001 From: wadii Date: Wed, 23 Sep 2026 16:27:49 +0200 Subject: [PATCH 6/6] test: drop comments that restate the line below them --- tests/sdk/events.test.ts | 6 ------ 1 file changed, 6 deletions(-) diff --git a/tests/sdk/events.test.ts b/tests/sdk/events.test.ts index 5cc950f..926fbd5 100644 --- a/tests/sdk/events.test.ts +++ b/tests/sdk/events.test.ts @@ -202,7 +202,6 @@ test('custom events are never deduplicated', async () => { }); test('flush posts the batch to the events endpoint', async () => { - // The endpoint is built from the events API URL whether or not it has a trailing slash. const processor = eventProcessor({ eventsApiUrl: 'http://testUrl' }); trackPurchase(processor); @@ -242,7 +241,6 @@ test('flush sends custom headers without letting them override the SDK headers', customHeaders: { 'X-Proxy-Token': 'secret', 'Flagsmith-SDK-User-Agent': 'not-the-sdk', - // Header names are case-insensitive: a variant would be merged, not overridden. 'x-environment-key': 'not-the-environment', 'USER-AGENT': 'not-the-sdk' } @@ -274,7 +272,6 @@ test('the buffer is flushed as soon as it reaches maxBuffer', async () => { trackPurchase(processor, 'user-456'); - // Posted by reaching maxBuffer, before anything asks for a flush. expect(fetch).toHaveBeenCalledTimes(1); expect(postedEvents()).toHaveLength(2); @@ -291,7 +288,6 @@ test('flush waits for a batch posted by the flush timer', async () => { processor.start(); trackPurchase(processor); - // The timer starts a flush that never settles until the response is resolved. await vi.advanceTimersByTimeAsync(10000); expect(fetch).toHaveBeenCalledTimes(1); @@ -342,7 +338,6 @@ test('flush does not wait for a batch started after it was called', async () => firstFlushed = true; }); - // Traffic keeps arriving while the first batch is on the wire. trackPurchase(processor, 'user-456'); let secondFlushed = false; const secondFlush = processor.flush().then(() => { @@ -373,7 +368,6 @@ test('flush waits for every in-flight batch even when one of them rejects', asyn }); const processor = eventProcessor({ maxBuffer: 1, logger }); - // Each event reaches maxBuffer and starts its own batch. trackPurchase(processor); trackPurchase(processor, 'user-456'); let flushed = false;