Skip to content
Merged
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
3 changes: 2 additions & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@ on:
- '*'
pull_request:
branches:
- '*'
- main
- development

jobs:
build:
Expand Down
2 changes: 1 addition & 1 deletion package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@splitsoftware/splitio-commons",
"version": "1.3.2-rc.4",
"version": "1.3.2-rc.5",
"description": "Split Javascript SDK common components",
"main": "cjs/index.js",
"module": "esm/index.js",
Expand Down
18 changes: 12 additions & 6 deletions src/storages/KeyBuilderSS.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,6 @@ export class KeyBuilderSS extends KeyBuilder {
return `${this.prefix}.segments.registered`;
}

private buildVersionablePrefix() {
return `${this.metadata.s}/${this.metadata.n}/${this.metadata.i}`;
}

buildImpressionsKey() {
return `${this.prefix}.impressions`;
}
Expand All @@ -35,6 +31,12 @@ export class KeyBuilderSS extends KeyBuilder {
return `${this.prefix}.events`;
}

searchPatternForSplitKeys() {
return `${this.buildSplitKeyPrefix()}*`;
}

/* Telemetry keys */

buildLatencyKey(method: Method, bucket: number) {
return `${this.prefix}.telemetry.latencies::${this.buildVersionablePrefix()}/${methodNames[method]}/${bucket}`;
}
Expand All @@ -43,8 +45,12 @@ export class KeyBuilderSS extends KeyBuilder {
return `${this.prefix}.telemetry.exceptions::${this.buildVersionablePrefix()}/${methodNames[method]}`;
}

searchPatternForSplitKeys() {
return `${this.buildSplitKeyPrefix()}*`;
buildInitKey() {
return `${this.prefix}.telemetry.init::${this.buildVersionablePrefix()}`;
}

private buildVersionablePrefix() {
return `${this.metadata.s}/${this.metadata.n}/${this.metadata.i}`;
}

}
2 changes: 1 addition & 1 deletion src/storages/inRedis/RedisAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ import { timeout } from '../../utils/promise/timeout';
const LOG_PREFIX = 'storage:redis-adapter: ';

// If we ever decide to fully wrap every method, there's a Commander.getBuiltinCommands from ioredis.
const METHODS_TO_PROMISE_WRAP = ['set', 'exec', 'del', 'get', 'keys', 'sadd', 'srem', 'sismember', 'smembers', 'incr', 'rpush', 'pipeline', 'expire', 'mget', 'lrange', 'ltrim'];
const METHODS_TO_PROMISE_WRAP = ['set', 'exec', 'del', 'get', 'keys', 'sadd', 'srem', 'sismember', 'smembers', 'incr', 'rpush', 'pipeline', 'expire', 'mget', 'lrange', 'ltrim', 'hset'];

// Not part of the settings since it'll vary on each storage. We should be removing storage specific logic from elsewhere.
const DEFAULT_OPTIONS = {
Expand Down
7 changes: 7 additions & 0 deletions src/storages/inRedis/TelemetryCacheInRedis.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@ import { KeyBuilderSS } from '../KeyBuilderSS';
import { ITelemetryCacheAsync } from '../types';
import { findLatencyIndex } from '../findLatencyIndex';
import { Redis } from 'ioredis';
import { getTelemetryConfigStats } from '../../sync/submitters/telemetrySubmitter';
import { CONSUMER_MODE, STORAGE_REDIS } from '../../utils/constants';

export class TelemetryCacheInRedis implements ITelemetryCacheAsync {

Expand All @@ -26,4 +28,9 @@ export class TelemetryCacheInRedis implements ITelemetryCacheAsync {
.catch(() => { /* Handle rejections for telemetry */ });
}

recordConfig() {
const [key, field] = this.keys.buildInitKey().split('::');
const value = JSON.stringify(getTelemetryConfigStats(CONSUMER_MODE, STORAGE_REDIS));
return this.redis.hset(key, field, value).catch(() => { /* Handle rejections for telemetry */ });
}
}
2 changes: 1 addition & 1 deletion src/storages/inRedis/__tests__/RedisAdapter.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ const LOG_PREFIX = 'storage:redis-adapter: ';
// Mocking ioredis

// The list of methods we're wrapping on a promise (for timeout) on the adapter.
const METHODS_TO_PROMISE_WRAP = ['set', 'exec', 'del', 'get', 'keys', 'sadd', 'srem', 'sismember', 'smembers', 'incr', 'rpush', 'pipeline', 'expire', 'mget'];
const METHODS_TO_PROMISE_WRAP = ['set', 'exec', 'del', 'get', 'keys', 'sadd', 'srem', 'sismember', 'smembers', 'incr', 'rpush', 'pipeline', 'expire', 'mget', 'lrange', 'ltrim', 'hset'];

const ioredisMock = reduce([...METHODS_TO_PROMISE_WRAP, 'disconnect'], (acc, methodName) => {
acc[methodName] = jest.fn(() => Promise.resolve(methodName));
Expand Down
15 changes: 14 additions & 1 deletion src/storages/inRedis/__tests__/TelemetryCacheInRedis.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,28 +7,41 @@ import { fakeMetadata } from '../../pluggable/__tests__/ImpressionsCachePluggabl
const prefix = 'telemetry_cache_ut';
const exceptionKey = `${prefix}.telemetry.exceptions`;
const latencyKey = `${prefix}.telemetry.latencies`;
const initKey = `${prefix}.telemetry.init`;
const fieldVersionablePrefix = `${fakeMetadata.s}/${fakeMetadata.n}/${fakeMetadata.i}`;

test('TELEMETRY CACHE IN REDIS / `recordLatency` and `recordException`', async () => {
test('TELEMETRY CACHE IN REDIS', async () => {

const keysBuilder = new KeyBuilderSS(prefix, fakeMetadata);
const connection = new Redis();
const cache = new TelemetryCacheInRedis(loggerMock, keysBuilder, connection);

// recordException
expect(await cache.recordException('tr')).toBe(1);
expect(await cache.recordException('tr')).toBe(2);

expect(await connection.hget(exceptionKey, fieldVersionablePrefix + '/track')).toBe('2');
expect(await connection.hget(exceptionKey, fieldVersionablePrefix + '/treatment')).toBe(null);

// recordLatency
expect(await cache.recordLatency('tr', 1.6)).toBe(1);
expect(await cache.recordLatency('tr', 1.6)).toBe(2);

expect(await connection.hget(latencyKey, fieldVersionablePrefix + '/track/2')).toBe('2');
expect(await connection.hget(latencyKey, fieldVersionablePrefix + '/treatment/2')).toBe(null);

// recordConfig
expect(await cache.recordConfig()).toBe(1);
expect(JSON.parse(await connection.hget(initKey, fieldVersionablePrefix) as string)).toEqual({
oM: 1,
st: 'redis',
aF: 0,
rF: 0
});

// Clean up then end.
await connection.hdel(exceptionKey, fieldVersionablePrefix + '/track');
await connection.hdel(latencyKey, fieldVersionablePrefix + '/track/2');
await connection.hdel(initKey, fieldVersionablePrefix);
await connection.quit();
});
6 changes: 5 additions & 1 deletion src/storages/inRedis/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,18 +26,22 @@ export function InRedisStorage(options: InRedisStorageOptions = {}): IStorageAsy

const keys = new KeyBuilderSS(prefix, metadata);
const redisClient = new RedisAdapter(log, options.options || {});
const telemetry = new TelemetryCacheInRedis(log, keys, redisClient);

// subscription to Redis connect event in order to emit SDK_READY event on consumer mode
redisClient.on('connect', () => {
onReadyCb();

// Synchronize config
telemetry.recordConfig();
});

return {
splits: new SplitsCacheInRedis(log, keys, redisClient),
segments: new SegmentsCacheInRedis(log, keys, redisClient),
impressions: new ImpressionsCacheInRedis(log, keys.buildImpressionsKey(), redisClient, metadata),
events: new EventsCacheInRedis(log, keys.buildEventsKey(), redisClient, metadata),
telemetry: new TelemetryCacheInRedis(log, keys, redisClient),
telemetry,

// When using REDIS we should:
// 1- Disconnect from the storage
Expand Down
3 changes: 2 additions & 1 deletion src/storages/pluggable/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ function validatePluggableStorageOptions(options: any) {
function wrapperConnect(wrapper: IPluggableStorageWrapper, onReadyCb: (error?: any) => void) {
wrapper.connect().then(() => {
onReadyCb();
// At the moment, we don't synchronize config with pluggable storage
}).catch((e) => {
onReadyCb(e || new Error('Error connecting wrapper'));
});
Expand Down Expand Up @@ -77,7 +78,7 @@ export function PluggableStorage(options: PluggableStorageOptions): IStorageAsyn
impressions: isPartialConsumer ? new ImpressionsCacheInMemory(impressionsQueueSize) : new ImpressionsCachePluggable(log, keys.buildImpressionsKey(), wrapper, metadata),
impressionCounts: optimize ? new ImpressionCountsCacheInMemory() : undefined,
events: isPartialConsumer ? promisifyEventsTrack(new EventsCacheInMemory(eventsQueueSize)) : new EventsCachePluggable(log, keys.buildEventsKey(), wrapper, metadata),
// @TODO Not using TelemetryCachePluggable yet, because it is not supported by the Split Synchronizer
// @TODO Not using TelemetryCachePluggable yet because it's not supported by the Split Synchronizer, and needs to drop or queue operations while the wrapper is not ready
// telemetry: isPartialConsumer ? new TelemetryCacheInMemory() : new TelemetryCachePluggable(log, keys, wrapper),

// Disconnect the underlying storage
Expand Down
2 changes: 1 addition & 1 deletion src/sync/submitters/__tests__/telemetrySubmitter.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,7 @@ describe('Telemetry submitter', () => {
expect(recordTimeUntilReadySpy).toBeCalledTimes(1);

expect(postMetricsConfig).toBeCalledWith(JSON.stringify({
oM: 0, st: 'memory', sE: true, rR: { sp: 1, se: 1, im: 1, ev: 1, te: 100 }, uO: { s: true, e: true, a: true, st: true, t: true }, iQ: 1, eQ: 1, iM: 0, iL: false, hP: false, aF: 0, rF: 0, tR: 0, tC: 0, nR: 0, t: [], i: ['NoopIntegration']
oM: 0, st: 'memory', aF: 0, rF: 0, sE: true, rR: { sp: 0.001, se: 0.001, im: 0.001, ev: 0.001, te: 0.1 }, uO: { s: true, e: true, a: true, st: true, t: true }, iQ: 1, eQ: 1, iM: 0, iL: false, hP: false, tR: 0, tC: 0, nR: 0, t: [], i: ['NoopIntegration'], uC: 0
}));

// Stop submitter, to not execute the 1st periodic metrics/usage POST
Expand Down
41 changes: 27 additions & 14 deletions src/sync/submitters/telemetrySubmitter.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,14 @@
import { ISegmentsCacheSync, ISplitsCacheSync, ITelemetryCacheSync } from '../../storages/types';
import { submitterFactory, firstPushWindowDecorator } from './submitter';
import { TelemetryUsageStatsPayload, TelemetryConfigStatsPayload } from './types';
import { QUEUED, DEDUPED, DROPPED, CONSUMER_MODE, CONSUMER_ENUM, STANDALONE_MODE, CONSUMER_PARTIAL_MODE, STANDALONE_ENUM, CONSUMER_PARTIAL_ENUM, OPTIMIZED, DEBUG, DEBUG_ENUM, OPTIMIZED_ENUM } from '../../utils/constants';
import { TelemetryUsageStatsPayload, TelemetryConfigStatsPayload, TelemetryConfigStats } from './types';
import { QUEUED, DEDUPED, DROPPED, CONSUMER_MODE, CONSUMER_ENUM, STANDALONE_MODE, CONSUMER_PARTIAL_MODE, STANDALONE_ENUM, CONSUMER_PARTIAL_ENUM, OPTIMIZED, DEBUG, DEBUG_ENUM, OPTIMIZED_ENUM, CONSENT_GRANTED, CONSENT_DECLINED, CONSENT_UNKNOWN } from '../../utils/constants';
import { SDK_READY, SDK_READY_FROM_CACHE } from '../../readiness/constants';
import { ISettings } from '../../types';
import { ConsentStatus, ISettings, SDKMode } from '../../types';
import { base } from '../../utils/settingsValidation';
import { usedKeysMap } from '../../utils/inputValidation/apiKey';
import { timer } from '../../utils/timeTracker/timer';
import { ISdkFactoryContextSync } from '../../sdkFactory/types';
import { objectAssign } from '../../utils/lang/objectAssign';

/**
* Converts data from telemetry cache into /metrics/usage request payload.
Expand Down Expand Up @@ -54,6 +55,12 @@ const IMPRESSIONS_MODE_MAP = {
[DEBUG]: DEBUG_ENUM
} as Record<ISettings['sync']['impressionsMode'], (0 | 1)>;

const USER_CONSENT_MAP = {
[CONSENT_UNKNOWN]: 1,
[CONSENT_GRANTED]: 2,
[CONSENT_DECLINED]: 3
} as Record<ConsentStatus, number>;

function getActiveFactories() {
return Object.keys(usedKeysMap).length;
}
Expand All @@ -64,6 +71,15 @@ function getRedundantActiveFactories() {
}, 0);
}

export function getTelemetryConfigStats(mode: SDKMode, storageType: string): TelemetryConfigStats {
return {
oM: OPERATION_MODE_MAP[mode], // @ts-ignore lower case of storage type
st: storageType.toLowerCase(),
aF: getActiveFactories(),
rF: getRedundantActiveFactories(),
};
}

/**
* Converts data from telemetry cache and settings into /metrics/config request payload.
*/
Expand All @@ -75,16 +91,14 @@ export function telemetryCacheConfigAdapter(telemetry: ITelemetryCacheSync, sett
state(): TelemetryConfigStatsPayload {
const { urls, scheduler } = settings;

return {
oM: OPERATION_MODE_MAP[settings.mode], // @ts-ignore lower case of storage type
st: settings.storage.type.toLowerCase(),
return objectAssign(getTelemetryConfigStats(settings.mode, settings.storage.type), {
sE: settings.streamingEnabled,
rR: {
sp: scheduler.featuresRefreshRate,
se: scheduler.segmentsRefreshRate,
im: scheduler.impressionsRefreshRate,
ev: scheduler.eventsPushRate,
te: scheduler.telemetryRefreshRate,
sp: scheduler.featuresRefreshRate / 1000,
se: scheduler.segmentsRefreshRate / 1000,
im: scheduler.impressionsRefreshRate / 1000,
ev: scheduler.eventsPushRate / 1000,
te: scheduler.telemetryRefreshRate / 1000,
}, // refreshRates
uO: {
s: urls.sdk !== base.urls.sdk,
Expand All @@ -98,14 +112,13 @@ export function telemetryCacheConfigAdapter(telemetry: ITelemetryCacheSync, sett
iM: IMPRESSIONS_MODE_MAP[settings.sync.impressionsMode],
iL: settings.impressionListener ? true : false,
hP: false, // @TODO proxy not supported
aF: getActiveFactories(),
rF: getRedundantActiveFactories(),
tR: telemetry.getTimeUntilReady() as number,
tC: telemetry.getTimeUntilReadyFromCache(),
nR: telemetry.getNonReadyUsage(),
t: telemetry.popTags(),
i: settings.integrations && settings.integrations.map(int => int.type),
};
uC: settings.userConsent ? USER_CONSENT_MAP[settings.userConsent] : 0
});
}
};
}
Expand Down
17 changes: 11 additions & 6 deletions src/sync/submitters/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -165,10 +165,17 @@ export type UrlOverrides = {
t: boolean, // telemetry
}

// 'metrics/config' JSON request body
export type TelemetryConfigStatsPayload = {
oM?: OperationMode, // operationMode
// 'telemetry.init' Redis/Pluggable key
export type TelemetryConfigStats = {
oM: OperationMode, // operationMode
st: 'memory' | 'redis' | 'pluggable' | 'localstorage', // storage
aF: number, // activeFactories
rF: number, // redundantActiveFactories
t?: Array<string>, // tags
}

// 'metrics/config' JSON request body
export type TelemetryConfigStatsPayload = TelemetryConfigStats & {
sE: boolean, // streamingEnabled
rR: RefreshRates, // refreshRates
uO: UrlOverrides, // urlOverrides
Expand All @@ -177,11 +184,9 @@ export type TelemetryConfigStatsPayload = {
iM: ImpressionsMode, // impressionsMode
iL: boolean, // impressionsListenerEnabled
hP: boolean, // httpProxyDetected
aF: number, // activeFactories
rF: number, // redundantActiveFactories
tR: number, // timeUntilSDKReady
tC?: number, // timeUntilSDKReadyFromCache
nR: number, // SDKNotReadyUsage
t?: Array<string>, // tags
i?: Array<string>, // integrations
uC: number, // userConsent
}
4 changes: 2 additions & 2 deletions src/trackers/telemetryTracker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ export function telemetryTrackerFactory(
return (error) => {
(telemetryCache as ITelemetryCacheSync).recordHttpLatency(operation, httpTime());
if (error && error.statusCode) (telemetryCache as ITelemetryCacheSync).recordHttpError(operation, error.statusCode);
else (telemetryCache as ITelemetryCacheSync).recordSuccessfulSync(operation, now());
else (telemetryCache as ITelemetryCacheSync).recordSuccessfulSync(operation, Date.now());
};
},
sessionLength() { // @ts-ignore ITelemetryCacheAsync doesn't implement the method
Expand All @@ -44,7 +44,7 @@ export function telemetryTrackerFactory(
(telemetryCache as ITelemetryCacheSync).recordAuthRejections();
} else {
(telemetryCache as ITelemetryCacheSync).recordStreamingEvents({
e, d, t: now()
e, d, t: Date.now()
});
if (e === TOKEN_REFRESH) (telemetryCache as ITelemetryCacheSync).recordTokenRefreshes();
}
Expand Down
2 changes: 1 addition & 1 deletion src/utils/timeTracker/now/node.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,5 +3,5 @@ export function now() {
// eslint-disable-next-line no-undef
let time = process.hrtime();

return time[0] * 1e3 + time[1] * 1e-6; // convert it to milis
return time[0] * 1e3 + time[1] * 1e-6; // convert it to millis
}