From 1bab7f7e7e5f53e0b30567c3970f1b14a2eb9329 Mon Sep 17 00:00:00 2001 From: Cathleen Yan <58714163+cathleeny@users.noreply.github.com> Date: Tue, 6 Oct 2026 21:29:37 +0000 Subject: [PATCH 1/2] Generalize driver feature flag cache and typed reads Signed-off-by: Cathleen Yan <58714163+cathleeny@users.noreply.github.com> --- lib/DBSQLClient.ts | 30 +++- lib/{telemetry => }/FeatureFlagCache.ts | 143 +++++++++++------ lib/telemetry/TelemetryClient.ts | 20 +-- lib/telemetry/TelemetryClientProvider.ts | 6 +- .../telemetry/telemetry-integration.test.ts | 6 +- tests/unit/DBSQLClient.test.ts | 13 +- tests/unit/telemetry/FeatureFlagCache.test.ts | 151 +++++++++++++----- tests/unit/telemetry/TelemetryClient.test.ts | 28 ---- 8 files changed, 244 insertions(+), 153 deletions(-) rename lib/{telemetry => }/FeatureFlagCache.ts (54%) diff --git a/lib/DBSQLClient.ts b/lib/DBSQLClient.ts index 02aa1d2e..c7b888f1 100644 --- a/lib/DBSQLClient.ts +++ b/lib/DBSQLClient.ts @@ -37,6 +37,7 @@ import CloseableCollection from './utils/CloseableCollection'; import IConnectionProvider from './connection/contracts/IConnectionProvider'; import TelemetryClient from './telemetry/TelemetryClient'; import TelemetryClientProvider from './telemetry/TelemetryClientProvider'; +import FeatureFlagCache from './FeatureFlagCache'; import TelemetryEventEmitter from './telemetry/TelemetryEventEmitter'; import MetricsAggregator from './telemetry/MetricsAggregator'; import { DriverConfiguration, DRIVER_NAME, TelemetryEventType, DEFAULT_TELEMETRY_CONFIG } from './telemetry/types'; @@ -118,6 +119,8 @@ export default class DBSQLClient extends EventEmitter implements IDBSQLClient, I private telemetryClient?: TelemetryClient; + private featureFlagCache?: FeatureFlagCache; + private telemetryEmitter?: TelemetryEventEmitter; // True once we've shipped the full DriverConfiguration on a CONNECTION_OPEN @@ -579,13 +582,15 @@ export default class DBSQLClient extends EventEmitter implements IDBSQLClient, I try { // Acquire (or create) the per-host TelemetryClient from the // process-wide provider. The shared client owns the circuit-breaker - // registry, feature-flag cache, exporter, and aggregator. Multiple + // registry, exporter, and aggregator. Multiple // DBSQLClient instances on the same host share these resources so // breaker counters and HTTP batches don't fragment per-instance. this.telemetryClient = TelemetryClientProvider.getInstance().getOrCreateClient(this, this.host); - // Use the shared feature-flag cache (registered in the previous step). - const enabled = await this.telemetryClient.getFeatureFlagCache().isTelemetryEnabled(this.host); + const enabled = await this.getFeatureFlagCache().getBoolean( + this.host, + 'databricks.partnerplatform.clientConfigsFeatureFlags.enableTelemetryForNodeJs', + ); if (!enabled) { // Release our refcount immediately; we won't be emitting. @@ -610,7 +615,7 @@ export default class DBSQLClient extends EventEmitter implements IDBSQLClient, I } catch (error: any) { // Swallow all telemetry initialization errors. If we acquired a refcount // before the throw, release it — otherwise the per-host TelemetryClient - // (and its flush timer / exporter / FFCache) leaks for the lifetime of + // (and its flush timer / exporter) leaks for the lifetime of // the process on long-running supervisors that retry-connect. if (this.telemetryClient) { try { @@ -628,6 +633,20 @@ export default class DBSQLClient extends EventEmitter implements IDBSQLClient, I } } + // Also usable after driver auth/transport setup, before selecting a backend. + private getFeatureFlagCache(): FeatureFlagCache { + if (!this.featureFlagCache) { + this.featureFlagCache = new FeatureFlagCache(this); + this.featureFlagCache.getOrCreateContext(this.host!); + } + return this.featureFlagCache; + } + + private releaseFeatureFlagCache(): void { + if (this.featureFlagCache && this.host) this.featureFlagCache.releaseContext(this.host); + this.featureFlagCache = undefined; + } + /** * Connects DBSQLClient to endpoint * @public @@ -661,6 +680,7 @@ export default class DBSQLClient extends EventEmitter implements IDBSQLClient, I this.telemetryClient = undefined; this.telemetryEmitter = undefined; } + this.releaseFeatureFlagCache(); // Re-arm: the new connection is a fresh client-config lineage even if // the host is the same. this.driverConfigShipped = false; @@ -911,6 +931,8 @@ export default class DBSQLClient extends EventEmitter implements IDBSQLClient, I // with client.close) cannot smuggle events into the closed aggregator. this.telemetryEmitter = undefined; + this.releaseFeatureFlagCache(); + this.client = undefined; this.connectionProvider = undefined; this.authProvider = undefined; diff --git a/lib/telemetry/FeatureFlagCache.ts b/lib/FeatureFlagCache.ts similarity index 54% rename from lib/telemetry/FeatureFlagCache.ts rename to lib/FeatureFlagCache.ts index e060bd93..d3e2ecbd 100644 --- a/lib/telemetry/FeatureFlagCache.ts +++ b/lib/FeatureFlagCache.ts @@ -15,33 +15,35 @@ */ import fetch, { RequestInit, Response, Request } from 'node-fetch'; -import IClientContext from '../contracts/IClientContext'; -import { LogLevel } from '../contracts/IDBSQLLogger'; -import IAuthentication from '../connection/contracts/IAuthentication'; -import { buildTelemetryUrl, normalizeHeaders } from './telemetryUtils'; -import buildUserAgentString from '../utils/buildUserAgentString'; -import driverVersion from '../version'; +import IClientContext from './contracts/IClientContext'; +import { LogLevel } from './contracts/IDBSQLLogger'; +import IAuthentication from './connection/contracts/IAuthentication'; +import { buildTelemetryUrl, normalizeHeaders } from './telemetry/telemetryUtils'; +import buildUserAgentString from './utils/buildUserAgentString'; +import driverVersion from './version'; export interface FeatureFlagContext { - telemetryEnabled?: boolean; + flags?: Map; + fetchPromise?: Promise; lastFetched?: Date; refCount: number; cacheDuration: number; } /** - * Per-host feature-flag cache used to gate telemetry emission. Responsibilities: + * Shared feature-flag values, independent of telemetry or an opened session. + * Acquire/release a context per consumer; each reader uses its caller's auth. + * Responsibilities: * - dedupe in-flight fetches (thundering-herd protection); * - ref-count so context goes away when the last consumer closes; * - clamp server-provided TTL into a safe band. * - * Shares HTTP plumbing (agent, user agent) with DatabricksTelemetryExporter. - * Consumer wiring lands in a later PR in this stack (see PR description). + * Workspace ID partitions SPOG traffic; normalized host is the fallback. */ export default class FeatureFlagCache { - private contexts: Map; + private static sharedContexts = new Map(); - private fetchPromises: Map> = new Map(); + private contexts = FeatureFlagCache.sharedContexts; private readonly userAgent: string; @@ -51,69 +53,112 @@ export default class FeatureFlagCache { private readonly MAX_CACHE_DURATION_S = 3600; - private readonly FEATURE_FLAG_NAME = 'databricks.partnerplatform.clientConfigsFeatureFlags.enableTelemetryForNodeJs'; - constructor(private context: IClientContext, private authProvider?: IAuthentication) { - this.contexts = new Map(); this.userAgent = buildUserAgentString(this.context.getConfig().userAgentEntry); } + private cacheKey(host: string): string { + const headers = this.context.getConfig().customHeaders ?? {}; + const workspace = Object.entries(headers).find(([name]) => name.toLowerCase() === 'x-databricks-org-id')?.[1]; + return workspace ? `workspace:${workspace}` : `host:${buildTelemetryUrl(host, '') ?? host}`; + } + getOrCreateContext(host: string): FeatureFlagContext { - let ctx = this.contexts.get(host); + const key = this.cacheKey(host); + let ctx = this.contexts.get(key); if (!ctx) { ctx = { refCount: 0, cacheDuration: this.CACHE_DURATION_MS, }; - this.contexts.set(host, ctx); + this.contexts.set(key, ctx); } ctx.refCount += 1; return ctx; } releaseContext(host: string): void { - const ctx = this.contexts.get(host); + const key = this.cacheKey(host); + const ctx = this.contexts.get(key); if (ctx) { ctx.refCount -= 1; if (ctx.refCount <= 0) { - this.contexts.delete(host); - this.fetchPromises.delete(host); + this.contexts.delete(key); } } } - async isTelemetryEnabled(host: string): Promise { + private async getRawValue(host: string, name: string): Promise { const logger = this.context.getLogger(); - const ctx = this.contexts.get(host); + const ctx = this.contexts.get(this.cacheKey(host)); if (!ctx) { - return false; + return undefined; } const isExpired = !ctx.lastFetched || Date.now() - ctx.lastFetched.getTime() > ctx.cacheDuration; if (isExpired) { - if (!this.fetchPromises.has(host)) { - const fetchPromise = this.fetchFeatureFlag(host) - .then((enabled) => { - ctx.telemetryEnabled = enabled; + if (!ctx.fetchPromise) { + ctx.fetchPromise = this.fetchFeatureFlags(host) + .then((flags) => { + ctx.flags = flags; ctx.lastFetched = new Date(); - return enabled; }) .catch((error: any) => { logger.log(LogLevel.debug, `Error fetching feature flag: ${error.message}`); - return ctx.telemetryEnabled ?? false; }) .finally(() => { - this.fetchPromises.delete(host); + ctx.fetchPromise = undefined; }); - this.fetchPromises.set(host, fetchPromise); } - await this.fetchPromises.get(host); + await ctx.fetchPromise; + } + + return ctx.flags?.get(name); + } + + private async getValue(host: string, name: string): Promise { + const raw = await this.getRawValue(host, name); + try { + return raw === undefined ? undefined : JSON.parse(raw); + } catch { + return undefined; } + } + + async getBoolean(host: string, name: string, defaultValue = false): Promise { + const value = await this.getValue(host, name); + return typeof value === 'boolean' ? value : defaultValue; + } + + async getInt32(host: string, name: string, defaultValue?: number): Promise { + const value = await this.getInt64(host, name); + return value !== undefined && value >= -2147483648 && value <= 2147483647 ? Number(value) : defaultValue; + } + + async getInt64(host: string, name: string, defaultValue?: bigint): Promise { + const raw = (await this.getRawValue(host, name))?.trim(); + // Parse integer text directly: JSON.parse would round values beyond 2^53. + if (!raw || !/^-?(0|[1-9]\d*)$/.test(raw)) return defaultValue; + const value = BigInt(raw); + return value >= BigInt('-9223372036854775808') && value <= BigInt('9223372036854775807') ? value : defaultValue; + } - return ctx.telemetryEnabled ?? false; + async getDouble(host: string, name: string, defaultValue?: number): Promise { + const value = await this.getValue(host, name); + return typeof value === 'number' && Number.isFinite(value) ? value : defaultValue; + } + + async getString(host: string, name: string, defaultValue?: string): Promise { + const value = await this.getValue(host, name); + return typeof value === 'string' ? value : defaultValue; + } + + async getStringList(host: string, name: string, defaultValue?: string[]): Promise { + const value = await this.getValue(host, name); + return Array.isArray(value) && value.every((item) => typeof item === 'string') ? value : defaultValue; } /** @@ -124,8 +169,9 @@ export default class FeatureFlagCache { return driverVersion.replace(/-oss$/, ''); } - private async fetchFeatureFlag(host: string): Promise { + private async fetchFeatureFlags(host: string): Promise> { const logger = this.context.getLogger(); + const ctx = this.contexts.get(this.cacheKey(host)); try { const endpoint = buildTelemetryUrl( @@ -134,7 +180,7 @@ export default class FeatureFlagCache { ); if (!endpoint) { logger.log(LogLevel.debug, `Feature flag fetch skipped: invalid host ${host}`); - return false; + return new Map(); } const headers: Record = { @@ -154,33 +200,30 @@ export default class FeatureFlagCache { if (!response.ok) { await response.text().catch(() => {}); - logger.log(LogLevel.debug, `Feature flag fetch failed: ${response.status} ${response.statusText}`); - return false; + throw new Error(`Feature flag fetch failed: ${response.status} ${response.statusText}`); } const data: any = await response.json(); if (data && data.flags && Array.isArray(data.flags)) { - const ctx = this.contexts.get(host); - if (ctx && typeof data.ttl_seconds === 'number' && data.ttl_seconds > 0) { + if (ctx && Number.isFinite(data.ttl_seconds) && data.ttl_seconds > 0) { const clampedTtl = Math.max(this.MIN_CACHE_DURATION_S, Math.min(this.MAX_CACHE_DURATION_S, data.ttl_seconds)); ctx.cacheDuration = clampedTtl * 1000; logger.log(LogLevel.debug, `Updated cache duration to ${clampedTtl} seconds`); } - const flag = data.flags.find((f: any) => f.name === this.FEATURE_FLAG_NAME); - if (flag) { - const enabled = String(flag.value).toLowerCase() === 'true'; - logger.log(LogLevel.debug, `Feature flag ${this.FEATURE_FLAG_NAME}: ${enabled}`); - return enabled; - } + return new Map( + data.flags + .filter((flag: any) => typeof flag?.name === 'string' && typeof flag.value === 'string') + .map((flag: any) => [flag.name, flag.value]), + ); } - logger.log(LogLevel.debug, `Feature flag ${this.FEATURE_FLAG_NAME} not found in response`); - return false; + return new Map(); } catch (error: any) { logger.log(LogLevel.debug, `Error fetching feature flag from ${host}: ${error.message}`); - return false; + // Preserve the existing failure policy: use defaults until the next TTL. + return new Map(); } } @@ -203,9 +246,7 @@ export default class FeatureFlagCache { } private async getAuthHeaders(): Promise> { - // Prefer the explicitly-injected auth provider; fall back to the context - // (used when a shared TelemetryClient resolves auth through its FIFO of - // registered DBSQLClients). Mirrors DatabricksTelemetryExporter.getAuthHeaders. + // Resolve auth from this caller, never from the shared cache state. const authProvider = this.authProvider ?? this.context.getAuthProvider?.(); if (!authProvider) { return {}; diff --git a/lib/telemetry/TelemetryClient.ts b/lib/telemetry/TelemetryClient.ts index 942a951e..c8f55e95 100644 --- a/lib/telemetry/TelemetryClient.ts +++ b/lib/telemetry/TelemetryClient.ts @@ -23,7 +23,6 @@ import IAuthentication from '../connection/contracts/IAuthentication'; import { CircuitBreakerRegistry, CircuitBreakerState } from './CircuitBreaker'; import DatabricksTelemetryExporter from './DatabricksTelemetryExporter'; import MetricsAggregator from './MetricsAggregator'; -import FeatureFlagCache from './FeatureFlagCache'; /** * Per-host telemetry resource owner. Held by `TelemetryClientProvider` @@ -31,7 +30,7 @@ import FeatureFlagCache from './FeatureFlagCache'; * connects to the same host. * * Owns the host-scoped triad — `MetricsAggregator`, `DatabricksTelemetryExporter`, - * `CircuitBreakerRegistry`, `FeatureFlagCache` — and implements `IClientContext` + * `CircuitBreakerRegistry` — and implements `IClientContext` * itself so those owned components have a stable context that outlives any * single `DBSQLClient`. The first registered `DBSQLClient`'s logger and config * are snapshotted; subsequent registrants donate their connection providers @@ -42,8 +41,6 @@ import FeatureFlagCache from './FeatureFlagCache'; * - Circuit-breaker state for `host` is correct only if all clients hitting * the same endpoint share counters (5 failures means 5 actual failures, not * 5×N for N independent `DBSQLClient` instances). - * - Feature-flag cache has a per-host TTL; deduping the GET prevents - * thundering-herd on cold cache. * - Metric batches mix events from every active client to the same host — * one HTTP POST per `flushIntervalMs` instead of N. */ @@ -56,8 +53,6 @@ class TelemetryClient implements IClientContext { private readonly circuitBreakerRegistry: CircuitBreakerRegistry; - private readonly featureFlagCache: FeatureFlagCache; - private readonly exporter: DatabricksTelemetryExporter; private readonly aggregator: MetricsAggregator; @@ -87,10 +82,6 @@ class TelemetryClient implements IClientContext { }); this.circuitBreakerRegistry = new CircuitBreakerRegistry(this); - this.featureFlagCache = new FeatureFlagCache(this); - // Register this host with the feature-flag cache so isTelemetryEnabled() - // does not short-circuit to false. close() releases via releaseContext(). - this.featureFlagCache.getOrCreateContext(host); this.exporter = new DatabricksTelemetryExporter(this, host, this.circuitBreakerRegistry); this.aggregator = new MetricsAggregator(this, this.exporter); @@ -283,10 +274,6 @@ class TelemetryClient implements IClientContext { return this.aggregator; } - getFeatureFlagCache(): FeatureFlagCache { - return this.featureFlagCache; - } - /** * Operator-visible snapshot of telemetry state for this host. Synchronous, * never throws — intended for health-check endpoints, shutdown banners, @@ -333,11 +320,6 @@ class TelemetryClient implements IClientContext { } catch (err) { this.logger.log(LogLevel.debug, `TelemetryClient exporter dispose error: ${(err as Error).message}`); } - try { - this.featureFlagCache.releaseContext(this.host); - } catch (err) { - this.logger.log(LogLevel.debug, `TelemetryClient FFCache release error: ${(err as Error).message}`); - } this.logger.log(LogLevel.debug, `Closed TelemetryClient for host: ${this.host}`); } } diff --git a/lib/telemetry/TelemetryClientProvider.ts b/lib/telemetry/TelemetryClientProvider.ts index eb873cde..e042a63b 100644 --- a/lib/telemetry/TelemetryClientProvider.ts +++ b/lib/telemetry/TelemetryClientProvider.ts @@ -31,8 +31,8 @@ const MAX_CLIENTS_SOFT_LIMIT = 128; /** * Process-wide registry of `TelemetryClient`s, one per host. Multiple * `DBSQLClient` instances connecting to the same host share the same - * `TelemetryClient`, which owns the host-scoped circuit breaker, feature - * flag cache, exporter, and aggregator. + * `TelemetryClient`, which owns the host-scoped circuit breaker, + * exporter, and aggregator. Feature flags have a separate workspace cache. * * Singleton because the resources we're sharing — circuit-breaker counters, * batched HTTP exports — are correct only at process scope. Per-`DBSQLClient` @@ -65,7 +65,7 @@ class TelemetryClientProvider { * Reset the process-wide singleton. Test-only — name-prefixed so * production callsites can't reach for it accidentally via autocomplete. * Resetting in production drops every host's circuit-breaker counters, - * feature-flag cache, exporter, and pending-metric buffer at once. + * exporter and pending-metric buffer at once. * * @internal Test-only. Production code MUST NOT call this. */ diff --git a/tests/e2e/telemetry/telemetry-integration.test.ts b/tests/e2e/telemetry/telemetry-integration.test.ts index 608378da..7888175f 100644 --- a/tests/e2e/telemetry/telemetry-integration.test.ts +++ b/tests/e2e/telemetry/telemetry-integration.test.ts @@ -17,7 +17,7 @@ import { expect } from 'chai'; import sinon from 'sinon'; import { DBSQLClient } from '../../../lib'; -import FeatureFlagCache from '../../../lib/telemetry/FeatureFlagCache'; +import FeatureFlagCache from '../../../lib/FeatureFlagCache'; import TelemetryClientProvider from '../../../lib/telemetry/TelemetryClientProvider'; import TelemetryEventEmitter from '../../../lib/telemetry/TelemetryEventEmitter'; import MetricsAggregator from '../../../lib/telemetry/MetricsAggregator'; @@ -163,7 +163,7 @@ describe('Telemetry Integration', () => { const client = new DBSQLClient(); // Stub feature flag to return false - const featureFlagStub = sinon.stub(FeatureFlagCache.prototype, 'isTelemetryEnabled').resolves(false); + const featureFlagStub = sinon.stub(FeatureFlagCache.prototype, 'getBoolean').resolves(false); try { await client.connect({ @@ -271,7 +271,7 @@ describe('Telemetry Integration', () => { // Stub feature flag to throw an error const featureFlagStub = sinon - .stub(FeatureFlagCache.prototype, 'isTelemetryEnabled') + .stub(FeatureFlagCache.prototype, 'getBoolean') .rejects(new Error('Feature flag fetch failed')); try { diff --git a/tests/unit/DBSQLClient.test.ts b/tests/unit/DBSQLClient.test.ts index 46f2709e..9996e165 100644 --- a/tests/unit/DBSQLClient.test.ts +++ b/tests/unit/DBSQLClient.test.ts @@ -22,7 +22,7 @@ import AuthProviderStub from './.stubs/AuthProviderStub'; import ConnectionProviderStub from './.stubs/ConnectionProviderStub'; import { TProtocolVersion } from '../../thrift/TCLIService_types'; import TelemetryClientProvider from '../../lib/telemetry/TelemetryClientProvider'; -import FeatureFlagCache from '../../lib/telemetry/FeatureFlagCache'; +import FeatureFlagCache from '../../lib/FeatureFlagCache'; import TelemetryEventEmitter from '../../lib/telemetry/TelemetryEventEmitter'; import { LogLevel } from '../../lib/contracts/IDBSQLLogger'; @@ -1127,8 +1127,9 @@ describe('DBSQLClient telemetry paths', () => { it('releases the prior refcount when connect() is called twice', async () => { const client = new DBSQLClient(); // Stub out feature-flag fetch to return true so the telemetry path runs. - sinon.stub(FeatureFlagCache.prototype, 'isTelemetryEnabled').resolves(true); + sinon.stub(FeatureFlagCache.prototype, 'getBoolean').resolves(true); const releaseSpy = sinon.spy(TelemetryClientProvider.prototype, 'releaseClient'); + const releaseFlagsSpy = sinon.spy(FeatureFlagCache.prototype, 'releaseContext'); await client.connect(connectOptions); // Sanity: refcount on host should be 1 after first connect. @@ -1137,17 +1138,19 @@ describe('DBSQLClient telemetry paths', () => { // Second connect to a different host should release the prior refcount. await client.connect({ ...connectOptions, host: '127.0.0.2' }); expect(releaseSpy.called, 'releaseClient should fire on reconnect').to.be.true; + expect(releaseFlagsSpy.calledWith(connectOptions.host)).to.be.true; // The old host should have decremented to 0 (closed and removed). expect(TelemetryClientProvider.getInstance().getRefCount(connectOptions.host)).to.equal(0); await client.close(); + expect(releaseFlagsSpy.callCount).to.equal(2); }); }); describe('telemetry refcount release path on init failure', () => { it('releases refcount when feature flag fetch throws', async () => { const client = new DBSQLClient(); - sinon.stub(FeatureFlagCache.prototype, 'isTelemetryEnabled').rejects(new Error('boom')); + sinon.stub(FeatureFlagCache.prototype, 'getBoolean').rejects(new Error('boom')); const releaseSpy = sinon.spy(TelemetryClientProvider.prototype, 'releaseClient'); await client.connect(connectOptions); @@ -1166,7 +1169,7 @@ describe('DBSQLClient telemetry paths', () => { const thriftClient = new ThriftClientStub(); sinon.stub(client, 'getClient').returns(Promise.resolve(thriftClient)); // Need a working emitter so we can spy on emitConnectionOpen. - sinon.stub(FeatureFlagCache.prototype, 'isTelemetryEnabled').resolves(true); + sinon.stub(FeatureFlagCache.prototype, 'getBoolean').resolves(true); const emitSpy = sinon.spy(TelemetryEventEmitter.prototype, 'emitConnectionOpen'); @@ -1196,7 +1199,7 @@ describe('DBSQLClient telemetry paths', () => { it('returns a populated snapshot when telemetry is enabled', async () => { const client = new DBSQLClient(); - sinon.stub(FeatureFlagCache.prototype, 'isTelemetryEnabled').resolves(true); + sinon.stub(FeatureFlagCache.prototype, 'getBoolean').resolves(true); await client.connect(connectOptions); const stats = client.getTelemetryStats(); expect(stats, 'stats should be populated when telemetry is on').to.not.be.undefined; diff --git a/tests/unit/telemetry/FeatureFlagCache.test.ts b/tests/unit/telemetry/FeatureFlagCache.test.ts index 85169a9f..8e212a32 100644 --- a/tests/unit/telemetry/FeatureFlagCache.test.ts +++ b/tests/unit/telemetry/FeatureFlagCache.test.ts @@ -16,7 +16,8 @@ import { expect } from 'chai'; import sinon from 'sinon'; -import FeatureFlagCache, { FeatureFlagContext } from '../../../lib/telemetry/FeatureFlagCache'; +import { Response } from 'node-fetch'; +import FeatureFlagCache from '../../../lib/FeatureFlagCache'; import ClientContextStub from '../.stubs/ClientContextStub'; import { LogLevel } from '../../../lib/contracts/IDBSQLLogger'; @@ -24,6 +25,7 @@ describe('FeatureFlagCache', () => { let clock: sinon.SinonFakeTimers; beforeEach(() => { + (FeatureFlagCache as any).sharedContexts.clear(); clock = sinon.useFakeTimers(); }); @@ -42,7 +44,7 @@ describe('FeatureFlagCache', () => { expect(ctx).to.not.be.undefined; expect(ctx.refCount).to.equal(1); expect(ctx.cacheDuration).to.equal(15 * 60 * 1000); // 15 minutes - expect(ctx.telemetryEnabled).to.be.undefined; + expect(ctx.flags).to.be.undefined; expect(ctx.lastFetched).to.be.undefined; }); @@ -123,12 +125,12 @@ describe('FeatureFlagCache', () => { }); }); - describe('isTelemetryEnabled', () => { + describe('getBoolean', () => { it('should return false for non-existent host', async () => { const context = new ClientContextStub(); const cache = new FeatureFlagCache(context); - const enabled = await cache.isTelemetryEnabled('non-existent-host.databricks.com'); + const enabled = await cache.getBoolean('non-existent-host.databricks.com', 'flag'); expect(enabled).to.be.false; }); @@ -137,11 +139,11 @@ describe('FeatureFlagCache', () => { const cache = new FeatureFlagCache(context); const host = 'test-host.databricks.com'; - // Stub the private fetchFeatureFlag method - const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlag').resolves(true); + // Stub the private fetchFeatureFlags method + const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlags').resolves(new Map([['flag', 'true']])); cache.getOrCreateContext(host); - const enabled = await cache.isTelemetryEnabled(host); + const enabled = await cache.getBoolean(host, 'flag'); expect(fetchStub.calledOnce).to.be.true; expect(fetchStub.calledWith(host)).to.be.true; @@ -155,19 +157,19 @@ describe('FeatureFlagCache', () => { const cache = new FeatureFlagCache(context); const host = 'test-host.databricks.com'; - const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlag').resolves(true); + const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlags').resolves(new Map([['flag', 'true']])); cache.getOrCreateContext(host); // First call - should fetch - await cache.isTelemetryEnabled(host); + await cache.getBoolean(host, 'flag'); expect(fetchStub.calledOnce).to.be.true; // Advance time by 10 minutes (less than 15 minute TTL) clock.tick(10 * 60 * 1000); // Second call - should use cached value - const enabled = await cache.isTelemetryEnabled(host); + const enabled = await cache.getBoolean(host, 'flag'); expect(fetchStub.calledOnce).to.be.true; // Still only called once expect(enabled).to.be.true; @@ -179,14 +181,14 @@ describe('FeatureFlagCache', () => { const cache = new FeatureFlagCache(context); const host = 'test-host.databricks.com'; - const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlag'); - fetchStub.onFirstCall().resolves(true); - fetchStub.onSecondCall().resolves(false); + const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlags'); + fetchStub.onFirstCall().resolves(new Map([['flag', 'true']])); + fetchStub.onSecondCall().resolves(new Map([['flag', 'false']])); cache.getOrCreateContext(host); // First call - should fetch - const enabled1 = await cache.isTelemetryEnabled(host); + const enabled1 = await cache.getBoolean(host, 'flag'); expect(enabled1).to.be.true; expect(fetchStub.calledOnce).to.be.true; @@ -194,7 +196,7 @@ describe('FeatureFlagCache', () => { clock.tick(16 * 60 * 1000); // Second call - should refetch due to expiration - const enabled2 = await cache.isTelemetryEnabled(host); + const enabled2 = await cache.getBoolean(host, 'flag'); expect(enabled2).to.be.false; expect(fetchStub.calledTwice).to.be.true; @@ -207,10 +209,10 @@ describe('FeatureFlagCache', () => { const cache = new FeatureFlagCache(context); const host = 'test-host.databricks.com'; - const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlag').rejects(new Error('Network error')); + const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlags').rejects(new Error('Network error')); cache.getOrCreateContext(host); - const enabled = await cache.isTelemetryEnabled(host); + const enabled = await cache.getBoolean(host, 'flag'); expect(enabled).to.be.false; expect(logSpy.calledWith(LogLevel.debug, 'Error fetching feature flag: Network error')).to.be.true; @@ -219,31 +221,31 @@ describe('FeatureFlagCache', () => { logSpy.restore(); }); - it('should not propagate exceptions from fetchFeatureFlag', async () => { + it('should not propagate exceptions from fetchFeatureFlags', async () => { const context = new ClientContextStub(); const cache = new FeatureFlagCache(context); const host = 'test-host.databricks.com'; - const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlag').rejects(new Error('Network error')); + const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlags').rejects(new Error('Network error')); cache.getOrCreateContext(host); // Should not throw - const enabled = await cache.isTelemetryEnabled(host); + const enabled = await cache.getBoolean(host, 'flag'); expect(enabled).to.equal(false); fetchStub.restore(); }); - it('should return false when telemetryEnabled is undefined', async () => { + it('should return false when the flag is missing', async () => { const context = new ClientContextStub(); const cache = new FeatureFlagCache(context); const host = 'test-host.databricks.com'; - const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlag').resolves(undefined); + const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlags').resolves(undefined); cache.getOrCreateContext(host); - const enabled = await cache.isTelemetryEnabled(host); + const enabled = await cache.getBoolean(host, 'flag'); expect(enabled).to.be.false; @@ -251,8 +253,8 @@ describe('FeatureFlagCache', () => { }); }); - describe('fetchFeatureFlag', () => { - it('should return false as placeholder implementation', async () => { + describe('fetchFeatureFlags', () => { + it('should default to false when the HTTP request fails', async () => { const context = new ClientContextStub(); const cache = new FeatureFlagCache(context); const host = 'test-host.databricks.com'; @@ -262,11 +264,11 @@ describe('FeatureFlagCache', () => { // (bogus) host; under mocha's 2s default this passed only when the // DNS failure happened to resolve quickly — flaky across runners / // Node versions (it timed out on Node 14/16/18 in CI). The behavior - // under test is just that `fetchFeatureFlag` resolves to `false`. + // under test is just that `fetchFeatureFlags` resolves to `false`. const fetchStub = sinon.stub(cache as any, 'fetchWithRetry').rejects(new Error('network disabled in test')); - // Access private method through any cast - const result = await (cache as any).fetchFeatureFlag(host); + cache.getOrCreateContext(host); + const result = await cache.getBoolean(host, 'flag'); expect(result).to.be.false; fetchStub.restore(); @@ -279,7 +281,7 @@ describe('FeatureFlagCache', () => { const cache = new FeatureFlagCache(context); const host = 'test-host.databricks.com'; - const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlag').resolves(true); + const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlags').resolves(new Map([['flag', 'true']])); // Simulate 3 connections to same host cache.getOrCreateContext(host); @@ -287,9 +289,9 @@ describe('FeatureFlagCache', () => { cache.getOrCreateContext(host); // All connections check telemetry - should only fetch once - await cache.isTelemetryEnabled(host); - await cache.isTelemetryEnabled(host); - await cache.isTelemetryEnabled(host); + await cache.getBoolean(host, 'flag'); + await cache.getBoolean(host, 'flag'); + await cache.getBoolean(host, 'flag'); expect(fetchStub.calledOnce).to.be.true; @@ -299,7 +301,7 @@ describe('FeatureFlagCache', () => { cache.releaseContext(host); // Context should be removed - const enabled = await cache.isTelemetryEnabled(host); + const enabled = await cache.getBoolean(host, 'flag'); expect(enabled).to.be.false; // No context, returns false fetchStub.restore(); @@ -311,15 +313,15 @@ describe('FeatureFlagCache', () => { const host1 = 'host1.databricks.com'; const host2 = 'host2.databricks.com'; - const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlag'); - fetchStub.withArgs(host1).resolves(true); - fetchStub.withArgs(host2).resolves(false); + const fetchStub = sinon.stub(cache as any, 'fetchFeatureFlags'); + fetchStub.withArgs(host1).resolves(new Map([['flag', 'true']])); + fetchStub.withArgs(host2).resolves(new Map([['flag', 'false']])); cache.getOrCreateContext(host1); cache.getOrCreateContext(host2); - const enabled1 = await cache.isTelemetryEnabled(host1); - const enabled2 = await cache.isTelemetryEnabled(host2); + const enabled1 = await cache.getBoolean(host1, 'flag'); + const enabled2 = await cache.getBoolean(host2, 'flag'); expect(enabled1).to.be.true; expect(enabled2).to.be.false; @@ -346,7 +348,7 @@ describe('FeatureFlagCache', () => { const cache = new FeatureFlagCache(context); const stub = sinon.stub(cache as any, 'fetchWithRetry').returns(makeJsonResponse({ flags: [] })); - await (cache as any).fetchFeatureFlag('host.example.com'); + await (cache as any).fetchFeatureFlags('host.example.com'); expect(stub.calledOnce).to.be.true; const init = stub.firstCall.args[1] as { headers: Record }; @@ -359,11 +361,80 @@ describe('FeatureFlagCache', () => { const cache = new FeatureFlagCache(context); const stub = sinon.stub(cache as any, 'fetchWithRetry').returns(makeJsonResponse({ flags: [] })); - await (cache as any).fetchFeatureFlag('host.example.com'); + await (cache as any).fetchFeatureFlags('host.example.com'); const init = stub.firstCall.args[1] as { headers: Record }; expect(init.headers).to.not.have.property('x-databricks-org-id'); stub.restore(); }); }); + + it('reads all six types from one GET without telemetry or a session', async () => { + const cases = [ + ['getBoolean', 'true', true], + ['getBoolean', '"true"', false], + ['getInt32', '2147483647', 2147483647], + ['getInt32', '2147483648', undefined], + ['getInt64', '9223372036854775807', BigInt('9223372036854775807')], + ['getInt64', '-9223372036854775808', BigInt('-9223372036854775808')], + ['getInt64', '9223372036854775808', undefined], + ['getInt64', '1.5', undefined], + ['getInt64', '01', undefined], + ['getInt64', 'true', undefined], + ['getDouble', '1.25', 1.25], + ['getDouble', '1e400', undefined], + ['getString', '"hello"', 'hello'], + ['getString', 'null', undefined], + ['getStringList', '["a","b"]', ['a', 'b']], + ['getStringList', '[null]', undefined], + ['getString', 'invalid', undefined], + ]; + const cache = new FeatureFlagCache(new ClientContextStub()); + cache.getOrCreateContext('test-host'); + const fetchStub = sinon.stub(cache as any, 'fetchWithRetry').resolves( + new Response( + JSON.stringify({ + flags: cases.map(([, value], index) => ({ name: String(index), value })), + ttl_seconds: 60, + }), + ), + ); + await Promise.all( + cases.map(async ([method, , expected], index) => { + expect(await (cache as any)[method as string]('test-host', String(index))).to.deep.equal(expected); + }), + ); + expect(await cache.getString('test-host', 'missing', 'fallback')).to.equal('fallback'); + expect(fetchStub.calledOnce).to.be.true; + }); + + it('shares values per workspace, refreshes with the current caller, and defaults on failure', async () => { + const first = new FeatureFlagCache(new ClientContextStub({ customHeaders: { 'x-databricks-org-id': '1' } })); + const current = new FeatureFlagCache(new ClientContextStub({ customHeaders: { 'X-Databricks-Org-Id': '1' } })); + const other = new FeatureFlagCache(new ClientContextStub({ customHeaders: { 'x-databricks-org-id': '2' } })); + expect(first.getOrCreateContext('host-a')).to.equal(current.getOrCreateContext('host-b')); + expect(other.getOrCreateContext('host-a')).to.not.equal(first.getOrCreateContext('host-a')); + const response = (value: string) => + new Response( + JSON.stringify({ + flags: [{ name: 'flag', value }], + ttl_seconds: 60, + }), + ); + const firstFetch = sinon.stub(first as any, 'fetchWithRetry').resolves(response('true')); + const currentFetch = sinon.stub(current as any, 'fetchWithRetry').resolves(response('false')); + expect(await Promise.all([first.getBoolean('host-a', 'flag'), current.getBoolean('host-b', 'flag')])).to.deep.equal( + [true, true], + ); + expect(currentFetch.called).to.be.false; + clock.tick(61000); + expect(await current.getBoolean('host-b', 'flag')).to.be.false; + expect(firstFetch.calledOnce).to.be.true; + expect(currentFetch.calledOnce).to.be.true; + clock.tick(61000); + currentFetch.rejects(new Error('offline')); + expect(await current.getBoolean('host-b', 'flag', true)).to.be.true; + expect(await current.getBoolean('host-b', 'flag')).to.be.false; + expect(currentFetch.calledTwice).to.be.true; + }); }); diff --git a/tests/unit/telemetry/TelemetryClient.test.ts b/tests/unit/telemetry/TelemetryClient.test.ts index ba8af8c7..ae05f0c9 100644 --- a/tests/unit/telemetry/TelemetryClient.test.ts +++ b/tests/unit/telemetry/TelemetryClient.test.ts @@ -150,34 +150,6 @@ describe('TelemetryClient', () => { }); }); - describe('feature-flag context registration (F1)', () => { - it('should register the host with the feature-flag cache on construction', () => { - const context = new ClientContextStub(); - const client = new TelemetryClient(context, HOST); - - // Without F1's wiring, isTelemetryEnabled returns false because the - // contexts map is empty. With F1, the constructor registers the host - // so the cache is ready to fetch the flag. - const cache = client.getFeatureFlagCache(); - // Internal access for assertion only — tests on getInstance/resetInstance - // would otherwise leak across the singleton. - const ctx = (cache as any).contexts.get(HOST); - expect(ctx, 'context should exist after TelemetryClient construction').to.exist; - expect(ctx.refCount, 'refCount should be 1 after registration').to.equal(1); - }); - - it('should release the feature-flag context on close', async () => { - const context = new ClientContextStub(); - const client = new TelemetryClient(context, HOST); - const cache = client.getFeatureFlagCache(); - - await client.close(); - - const ctx = (cache as any).contexts.get(HOST); - expect(ctx, 'context should be removed on close (refCount → 0)').to.be.undefined; - }); - }); - describe('multi-context FIFO', () => { it('registerContext appends contexts in registration order', () => { const ctxA = new ClientContextStub(); From aa104c9d58b61cd7fe76aa3ac725854f7d58917f Mon Sep 17 00:00:00 2001 From: Cathleen Yan <58714163+cathleeny@users.noreply.github.com> Date: Tue, 6 Oct 2026 22:44:17 +0000 Subject: [PATCH 2/2] refactor(feature-flags): simplify numeric validation Signed-off-by: Cathleen Yan <58714163+cathleeny@users.noreply.github.com> --- lib/FeatureFlagCache.ts | 13 ++++--------- 1 file changed, 4 insertions(+), 9 deletions(-) diff --git a/lib/FeatureFlagCache.ts b/lib/FeatureFlagCache.ts index d3e2ecbd..9ba00c8c 100644 --- a/lib/FeatureFlagCache.ts +++ b/lib/FeatureFlagCache.ts @@ -31,7 +31,7 @@ export interface FeatureFlagContext { } /** - * Shared feature-flag values, independent of telemetry or an opened session. + * Shared feature-flag values per workspace. * Acquire/release a context per consumer; each reader uses its caller's auth. * Responsibilities: * - dedupe in-flight fetches (thundering-herd protection); @@ -135,7 +135,7 @@ export default class FeatureFlagCache { async getInt32(host: string, name: string, defaultValue?: number): Promise { const value = await this.getInt64(host, name); - return value !== undefined && value >= -2147483648 && value <= 2147483647 ? Number(value) : defaultValue; + return value !== undefined && BigInt.asIntN(32, value) === value ? Number(value) : defaultValue; } async getInt64(host: string, name: string, defaultValue?: bigint): Promise { @@ -143,7 +143,7 @@ export default class FeatureFlagCache { // Parse integer text directly: JSON.parse would round values beyond 2^53. if (!raw || !/^-?(0|[1-9]\d*)$/.test(raw)) return defaultValue; const value = BigInt(raw); - return value >= BigInt('-9223372036854775808') && value <= BigInt('9223372036854775807') ? value : defaultValue; + return BigInt.asIntN(64, value) === value ? value : defaultValue; } async getDouble(host: string, name: string, defaultValue?: number): Promise { @@ -227,12 +227,7 @@ export default class FeatureFlagCache { } } - /** - * Retries transient network errors once before giving up. Without a retry - * a single hiccup would leave telemetry disabled for the full cache TTL - * (15 min). One retry gives an ephemeral DNS / connection-reset failure - * a second chance without pushing sustained load at a broken endpoint. - */ + /** Retries transient failures using the connection's retry policy. */ private async fetchWithRetry(url: string, init: RequestInit): Promise { const connectionProvider = await this.context.getConnectionProvider(); const agent = await connectionProvider.getAgent();