Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 26 additions & 4 deletions lib/DBSQLClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand All @@ -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 {
Expand All @@ -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
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down
150 changes: 93 additions & 57 deletions lib/telemetry/FeatureFlagCache.ts → lib/FeatureFlagCache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, string>;
fetchPromise?: Promise<void>;
lastFetched?: Date;
refCount: number;
cacheDuration: number;
}

/**
* Per-host feature-flag cache used to gate telemetry emission. Responsibilities:
* 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);
* - 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<string, FeatureFlagContext>;
private static sharedContexts = new Map<string, FeatureFlagContext>();

private fetchPromises: Map<string, Promise<boolean>> = new Map();
private contexts = FeatureFlagCache.sharedContexts;

private readonly userAgent: string;

Expand All @@ -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<boolean> {
private async getRawValue(host: string, name: string): Promise<string | undefined> {
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.telemetryEnabled ?? false;
return ctx.flags?.get(name);
}

private async getValue(host: string, name: string): Promise<unknown> {
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<boolean> {
const value = await this.getValue(host, name);
return typeof value === 'boolean' ? value : defaultValue;
}

async getInt32(host: string, name: string, defaultValue?: number): Promise<number | undefined> {
const value = await this.getInt64(host, name);
return value !== undefined && BigInt.asIntN(32, value) === value ? Number(value) : defaultValue;
}

async getInt64(host: string, name: string, defaultValue?: bigint): Promise<bigint | undefined> {
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 BigInt.asIntN(64, value) === value ? value : defaultValue;
}

async getDouble(host: string, name: string, defaultValue?: number): Promise<number | undefined> {
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<string | undefined> {
const value = await this.getValue(host, name);
return typeof value === 'string' ? value : defaultValue;
}

async getStringList(host: string, name: string, defaultValue?: string[]): Promise<string[] | undefined> {
const value = await this.getValue(host, name);
return Array.isArray(value) && value.every((item) => typeof item === 'string') ? value : defaultValue;
}

/**
Expand All @@ -124,8 +169,9 @@ export default class FeatureFlagCache {
return driverVersion.replace(/-oss$/, '');
}

private async fetchFeatureFlag(host: string): Promise<boolean> {
private async fetchFeatureFlags(host: string): Promise<Map<string, string>> {
const logger = this.context.getLogger();
const ctx = this.contexts.get(this.cacheKey(host));

try {
const endpoint = buildTelemetryUrl(
Expand All @@ -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<string, string> = {
Expand All @@ -154,42 +200,34 @@ 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();
}
}

/**
* 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<Response> {
const connectionProvider = await this.context.getConnectionProvider();
const agent = await connectionProvider.getAgent();
Expand All @@ -203,9 +241,7 @@ export default class FeatureFlagCache {
}

private async getAuthHeaders(): Promise<Record<string, string>> {
// 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 {};
Expand Down
Loading
Loading