Skip to content
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@ to docs, or any other relevant information.

### Added

- **Experimental**: Added `list` and `fetchHistory` interception to `WorkflowClientInterceptor`.
`list` interception spans one lazy iterable consumption across pagination and early termination.
- **Experimental**: Workflow Clients can now use `TypeInfo` to encode Workflow inputs and decode Workflow results.
- **Experimental**: `@temporalio/google-adk-agents` package for running Google ADK agents as durable Temporal Workflows.
ADK's OpenTelemetry agent-loop spans can be exported replay-safely from the Workflow sandbox by composing with
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,189 @@
import test from 'ava';
import * as otel from '@opentelemetry/api';
import { SpanStatusCode } from '@opentelemetry/api';
import { BasicTracerProvider, InMemorySpanExporter, SimpleSpanProcessor } from '@opentelemetry/sdk-trace-base';
import { instrumentAsyncIterable } from '../instrumentation';

/**
* Minimal sync ContextManager so unit tests can assert active-context propagation
* without adding @opentelemetry/context-async-hooks as a dependency.
*/
class TestContextManager implements otel.ContextManager {
private current: otel.Context = otel.ROOT_CONTEXT;

active(): otel.Context {
return this.current;
}

with<A extends unknown[], F extends (...args: A) => ReturnType<F>>(
context: otel.Context,
fn: F,
thisArg?: ThisParameterType<F>,
...args: A
): ReturnType<F> {
const previous = this.current;
this.current = context;
try {
return Reflect.apply(fn, thisArg, args);
} finally {
this.current = previous;
}
}

bind<T>(context: otel.Context, target: T): T {
if (typeof target !== 'function') {
return target;
}
// eslint-disable-next-line @typescript-eslint/no-this-alias -- bind() needs the manager instance in the wrapper
const manager = this;
const bound = function (this: unknown, ...args: unknown[]) {
return manager.with(context, () => (target as (...args: unknown[]) => unknown).apply(this, args));
};
return bound as T;
}

enable(): this {
return this;
}

disable(): this {
this.current = otel.ROOT_CONTEXT;
return this;
}
}

function setupTracer(name: string) {
const memoryExporter = new InMemorySpanExporter();
const provider = new BasicTracerProvider({
spanProcessors: [new SimpleSpanProcessor(memoryExporter)],
});
otel.context.setGlobalContextManager(new TestContextManager().enable());
return {
memoryExporter,
tracer: provider.getTracer(name),
};
}

test.afterEach.always(() => {
otel.trace.disable();
otel.context.disable();
});

async function* values(items: number[]): AsyncIterable<number> {
for (const item of items) {
yield item;
}
}

test('instrumentAsyncIterable does not open a span until iteration begins', async (t) => {
const { memoryExporter, tracer } = setupTracer('lazy-start');
const iterable = instrumentAsyncIterable({
tracer,
spanName: 'ListWorkflows',
fn: () => values([1, 2, 3]),
});

t.is(memoryExporter.getFinishedSpans().length, 0);
const iterator = iterable[Symbol.asyncIterator]();
t.is(memoryExporter.getFinishedSpans().length, 0);

t.is((await iterator.next()).value, 1);
t.is(memoryExporter.getFinishedSpans().length, 0);

t.is((await iterator.next()).value, 2);
t.is((await iterator.next()).value, 3);
t.true((await iterator.next()).done);

const spans = memoryExporter.getFinishedSpans();
t.is(spans.length, 1);
t.is(spans[0]!.name, 'ListWorkflows');
t.is(spans[0]!.status.code, SpanStatusCode.OK);
});

test('instrumentAsyncIterable ends the span exactly once on early break', async (t) => {
const { memoryExporter, tracer } = setupTracer('early-break');
const iterable = instrumentAsyncIterable({
tracer,
spanName: 'ListWorkflows',
fn: () => values([1, 2, 3, 4]),
});

const seen: number[] = [];
for await (const item of iterable) {
seen.push(item);
if (item === 2) {
t.is(memoryExporter.getFinishedSpans().length, 0);
break;
}
}

t.deepEqual(seen, [1, 2]);
const spans = memoryExporter.getFinishedSpans();
t.is(spans.length, 1);
t.is(spans[0]!.status.code, SpanStatusCode.OK);
});

test('instrumentAsyncIterable records errors and ends once', async (t) => {
const { memoryExporter, tracer } = setupTracer('error');
const error = new Error('downstream failed');
async function* boom(): AsyncIterable<number> {
yield 1;
throw error;
}

const iterable = instrumentAsyncIterable({
tracer,
spanName: 'ListWorkflows',
fn: () => boom(),
});

await t.throwsAsync(
async () => {
for await (const _ of iterable) {
// consume until error
}
},
{ is: error }
);

const spans = memoryExporter.getFinishedSpans();
t.is(spans.length, 1);
t.is(spans[0]!.status.code, SpanStatusCode.ERROR);
t.is(spans[0]!.status.message, error.message);
t.is(spans[0]!.events.filter((event) => event.name === 'exception').length, 1);
});

test('instrumentAsyncIterable propagates active context to downstream next()', async (t) => {
const { memoryExporter, tracer } = setupTracer('context');
const iterable = instrumentAsyncIterable({
tracer,
spanName: 'ListWorkflows',
fn: () => ({
[Symbol.asyncIterator]() {
let done = false;
return {
async next(): Promise<IteratorResult<number>> {
if (done) {
return { done: true, value: undefined };
}
done = true;
const child = tracer.startSpan('child-during-next');
child.end();
return { done: false, value: 1 };
},
};
},
}),
});

for await (const _ of iterable) {
// consume
}

const spans = memoryExporter.getFinishedSpans();
const parent = spans.find((span) => span.name === 'ListWorkflows');
const child = spans.find((span) => span.name === 'child-during-next');
t.truthy(parent);
t.truthy(child);
t.is(child!.parentSpanContext?.spanId, parent!.spanContext().spanId);
});
30 changes: 30 additions & 0 deletions contrib/interceptors-opentelemetry-v2/src/client/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,13 +13,18 @@ import type {
WorkflowTerminateInput,
WorkflowCancelInput,
WorkflowDescribeInput,
WorkflowFetchHistoryInput,
WorkflowListInput,
WorkflowClientInterceptor,
TerminateWorkflowExecutionResponse,
RequestCancelWorkflowExecutionResponse,
DescribeWorkflowExecutionResponse,
WorkflowExecutionInfo,
} from '@temporalio/client';
import type { History } from '@temporalio/common/lib/proto-utils';
import {
instrument,
instrumentAsyncIterable,
headersWithContext,
RUN_ID_ATTR_KEY,
WORKFLOW_ID_ATTR_KEY,
Expand Down Expand Up @@ -217,4 +222,29 @@ export class OpenTelemetryWorkflowClientInterceptor implements WorkflowClientInt
},
});
}

async fetchHistory(
input: WorkflowFetchHistoryInput,
next: Next<WorkflowClientInterceptor, 'fetchHistory'>
): Promise<History> {
return await instrument({
tracer: this.tracer,
spanName: SpanName.WORKFLOW_FETCH_HISTORY,
fn: async (span) => {
span.setAttribute(WORKFLOW_ID_ATTR_KEY, input.workflowExecution.workflowId);
if (input.workflowExecution.runId) {
span.setAttribute(RUN_ID_ATTR_KEY, input.workflowExecution.runId);
}
return await next(input);
},
});
}

list(input: WorkflowListInput, next: Next<WorkflowClientInterceptor, 'list'>): AsyncIterable<WorkflowExecutionInfo> {
return instrumentAsyncIterable({
tracer: this.tracer,
spanName: SpanName.WORKFLOW_LIST,
fn: () => next(input),
});
}
}
100 changes: 100 additions & 0 deletions contrib/interceptors-opentelemetry-v2/src/instrumentation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,10 @@ export interface InstrumentOptions<T> {

export type InstrumentOptionsSync<T> = Omit<InstrumentOptions<T>, 'fn'> & { fn: (span: otel.Span) => T };

export type InstrumentOptionsAsyncIterable<T> = Omit<InstrumentOptions<T>, 'fn'> & {
fn: (span: otel.Span) => AsyncIterable<T>;
};

/**
* Wraps `fn` in a span which ends when function returns or throws
*/
Expand Down Expand Up @@ -157,3 +161,99 @@ export function instrumentSync<T>({ tracer, spanName, fn, context, acceptableErr
}
return tracer.startActiveSpan(spanName, (span) => wrapWithSpanSync(span, fn, acceptableErrors));
}

/**
* Wraps an async iterable in a span whose lifetime matches iteration.
*
* The returned iterable is lazy: creating it does not open a span. The span starts on the first
* `next()` call, remains open while the iterator is in use, and ends on completion, early
* termination (`return`), or error.
*/
export function instrumentAsyncIterable<T>({
tracer,
spanName,
fn,
context,
acceptableErrors,
}: InstrumentOptionsAsyncIterable<T>): AsyncIterable<T> {
return {
[Symbol.asyncIterator](): AsyncIterator<T> {
let span: otel.Span | undefined;
let spanContext: otel.Context | undefined;
let iterator: AsyncIterator<T> | undefined;
let finished = false;

const finish = (err?: unknown): void => {
if (finished || span === undefined) {
return;
}
finished = true;
if (err !== undefined) {
maybeAddErrorToSpan(err, span, acceptableErrors);
} else {
span.setStatus({ code: otel.SpanStatusCode.OK });
}
span.end();
};

const ensureStarted = (): void => {
if (iterator !== undefined) {
return;
}
const parentContext = context ?? otel.context.active();
span = tracer.startSpan(spanName, undefined, parentContext);
spanContext = otel.trace.setSpan(parentContext, span);
const iterable = otel.context.with(spanContext, () => fn(span!));
iterator = iterable[Symbol.asyncIterator]();
};

return {
async next(...args: [] | [undefined]): Promise<IteratorResult<T>> {
try {
ensureStarted();
const result = await otel.context.with(spanContext!, async () => iterator!.next(...args));
if (result.done) {
finish();
}
return result;
} catch (err) {
finish(err);
throw err;
}
},
async return(value?: unknown): Promise<IteratorResult<T>> {
if (iterator === undefined) {
return { done: true, value: undefined as any };
}
try {
const result = iterator.return
? await otel.context.with(spanContext!, async () => iterator!.return!(value))
: ({ done: true, value: undefined } as IteratorResult<T>);
finish();
return result;
} catch (err) {
finish(err);
throw err;
}
},
async throw(err?: unknown): Promise<IteratorResult<T>> {
try {
ensureStarted();
if (iterator?.throw) {
const result = await otel.context.with(spanContext!, async () => iterator!.throw!(err));
if (result.done) {
finish();
}
return result;
}
finish(err);
throw err;
} catch (e) {
finish(e);
throw e;
}
},
};
},
};
}
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,16 @@ export enum SpanName {
*/
WORKFLOW_DESCRIBE = 'DescribeWorkflow',

/**
* Workflow history is fetched
*/
WORKFLOW_FETCH_HISTORY = 'FetchWorkflowHistory',

/**
* Workflows are listed
*/
WORKFLOW_LIST = 'ListWorkflows',

/**
* Workflow run is executing
*/
Expand Down
Loading