Skip to content
Merged
48 changes: 48 additions & 0 deletions packages/opentelemetry-tracing/src/BasicTracerProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,54 @@ export class BasicTracerProvider implements TracerProvider {
}
}

forceFlush(): Promise<void> {
const timeout = this._config.forceFlushTimeoutMillis;
const promises = this._registeredSpanProcessors.map(
(spanProcessor: SpanProcessor) => {
return new Promise(resolve => {
let state: 'resolved' | 'timeout' | 'error' | 'unresolved' =
Comment thread
weyert marked this conversation as resolved.
Outdated
'unresolved';
const timeoutInterval = setTimeout(() => {
resolve(
new Error(
`Span processor did not completed within timeout period of ${timeout} ms`
)
);
state = 'timeout';
}, timeout);

spanProcessor
.forceFlush()
.then(() => {
clearTimeout(timeoutInterval);
if (state !== 'timeout') {
state = 'resolved';
resolve(state);
}
})
.catch(error => {
clearTimeout(timeoutInterval);
state = 'error';
resolve(error);
});
});
}
);

return new Promise<void>((resolve, reject) => {
Promise.all(promises)
.then(results => {
const errors = results.filter(result => result !== 'resolved');
if (errors.length > 0) {
reject(errors);
} else {
resolve();
}
})
.catch(error => reject([error]));
});
}

shutdown() {
return this.activeSpanProcessor.shutdown();
}
Expand Down
1 change: 1 addition & 0 deletions packages/opentelemetry-tracing/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ const FALLBACK_OTEL_TRACES_SAMPLER = TracesSamplerValues.AlwaysOn;
*/
export const DEFAULT_CONFIG = {
sampler: buildSamplerFromEnv(env),
forceFlushTimeoutMillis: 30000,
traceParams: {
numberOfAttributesPerSpan: getEnv().OTEL_SPAN_ATTRIBUTE_COUNT_LIMIT,
numberOfLinksPerSpan: getEnv().OTEL_SPAN_LINK_COUNT_LIMIT,
Expand Down
6 changes: 6 additions & 0 deletions packages/opentelemetry-tracing/src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,12 @@ export interface TracerConfig {
* The default idGenerator generates random ids
*/
idGenerator?: IdGenerator;

/**
* How long the forceFlush can run before it is cancelled.
* The default value is 30000ms
*/
forceFlushTimeoutMillis?: number;
}

/**
Expand Down
50 changes: 50 additions & 0 deletions packages/opentelemetry-tracing/test/BasicTracerProvider.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ import { Resource } from '@opentelemetry/resources';
import * as assert from 'assert';
import * as sinon from 'sinon';
import { BasicTracerProvider, Span } from '../src';
import { NoopSpanProcessor } from '../src/NoopSpanProcessor';

describe('BasicTracerProvider', () => {
let removeEvent: Function | undefined;
Expand Down Expand Up @@ -392,6 +393,55 @@ describe('BasicTracerProvider', () => {
});
});

describe('.forceFlush()', () => {
it('should call forceFlush on all registered span processors', done => {
const tracerProvider = new BasicTracerProvider();
const spanProcessorOne = new NoopSpanProcessor();
const spyOne = sinon.spy(spanProcessorOne, 'forceFlush');
const spanProcessorTwo = new NoopSpanProcessor();
const spyTwo = sinon.spy(spanProcessorTwo, 'forceFlush');

tracerProvider.addSpanProcessor(spanProcessorOne);
tracerProvider.addSpanProcessor(spanProcessorTwo);

tracerProvider
.forceFlush()
.then(() => {
assert(spyOne.calledOnce);
assert(spyTwo.calledOnce);
done();
})
.catch(error => {
done(error);
});
});

it('should throw error when calling forceFlush on all registered span processors fails', done => {
const forceFlushStub = sinon.stub(
NoopSpanProcessor.prototype,
'forceFlush'
);
forceFlushStub.returns(Promise.reject('Error'));

const tracerProvider = new BasicTracerProvider();
const spanProcessorOne = new NoopSpanProcessor();
const spanProcessorTwo = new NoopSpanProcessor();
tracerProvider.addSpanProcessor(spanProcessorOne);
tracerProvider.addSpanProcessor(spanProcessorTwo);

tracerProvider
.forceFlush()
.then(() => {
done(new Error('Successful forceFlush not expected'));
})
.catch(_error => {
forceFlushStub.restore();
sinon.assert.calledTwice(forceFlushStub);
done();
});
});
});

describe('.bind()', () => {
it('should bind context with NoopContextManager context manager', done => {
const tracer = new BasicTracerProvider().getTracer('default');
Expand Down