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
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
import { BaseCallbackHandler } from '@langchain/core/callbacks/base';
import { RunnableLambda } from '@langchain/core/runnables';
import { TypeSafeClassifier } from '@langchain/typesafe';
// Sets up the AsyncLocalStorage that passes a parent run's config to its children, as it is inside `createAgent`.
import '@langchain/langgraph';
import * as Sentry from '@sentry/node';
import express from 'express';

function startMockTypeSafeServer() {
const app = express();
app.use(express.json());

app.post('/v1/systemone', (req, res) => {
res.json({
model: 'jev-1.13',
answers: { urgent: { type: 'noul', noul: 0.9 } },
usage: { input_tokens: 30, output_tokens: 2 },
});
});

return new Promise(resolve => {
const server = app.listen(0, () => {
resolve(server);
});
});
}

// Stands in for a tracer such as LangSmith, which the user passes to the parent run.
class RecordingHandler extends BaseCallbackHandler {
name = 'RecordingHandler';
runs = [];

handleChainStart(chain, _inputs, _runId, parentRunId) {
this.runs.push(`${chain.id.at(-1)}:${parentRunId ? 'child' : 'root'}`);
}
}

async function run() {
const server = await startMockTypeSafeServer();
const baseUrl = `http://localhost:${server.address().port}`;

await Sentry.startSpan({ op: 'function', name: 'main' }, async span => {
const classifier = new TypeSafeClassifier({
apiKey: 'mock-api-key',
baseUrl,
questions: { urgent: { type: 'noul', instructions: 'Is this urgent?' } },
});

// Like the `@langchain/typesafe` middlewares, call the classifier without a config.
const parent = RunnableLambda.from(input => classifier.invoke(input)).withConfig({ runName: 'parent' });
const recorder = new RecordingHandler();
await parent.invoke('My payouts have been failing.', { callbacks: [recorder] });

span.setAttribute('test.recorded_runs', recorder.runs.join(','));
});

await Sentry.flush(2000);
server.close();
}

run();
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
import { HumanMessage } from '@langchain/core/messages';
import { TypeSafeClassifier } from '@langchain/typesafe';
import * as Sentry from '@sentry/node';
import express from 'express';

function startMockTypeSafeServer() {
const app = express();
app.use(express.json());

app.post('/v1/systemone', (req, res) => {
res.json({
model: 'jev-1.13',
answers: { urgent: { type: 'noul', noul: 0.9 } },
usage: { input_tokens: 30, output_tokens: 2 },
});
});

return new Promise(resolve => {
const server = app.listen(0, () => {
resolve(server);
});
});
}

async function run() {
const server = await startMockTypeSafeServer();
const baseUrl = `http://localhost:${server.address().port}`;

await Sentry.startSpan({ op: 'function', name: 'main' }, async () => {
const classifier = new TypeSafeClassifier({
apiKey: 'mock-api-key',
baseUrl,
questions: { urgent: { type: 'noul', instructions: 'Is this urgent?' } },
});

await classifier.invoke('My payouts have been failing.');
await classifier.invoke(new HumanMessage('My card was charged twice.'));
});

server.close();
}

run();
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import {
GEN_AI_CONVERSATION_ID,
GEN_AI_INPUT_MESSAGES,
GEN_AI_OPERATION_NAME,
GEN_AI_OUTPUT_MESSAGES,
GEN_AI_PROVIDER_NAME,
GEN_AI_REQUEST_MAX_TOKENS,
GEN_AI_REQUEST_MODEL,
Expand All @@ -18,6 +19,7 @@ import {
SENTRY_OP,
SENTRY_ORIGIN,
} from '@sentry/conventions/attributes';
import { GEN_AI_EVALUATE } from '@sentry/conventions/op';
import { GEN_AI_RESPONSE_STOP_REASON_ATTRIBUTE } from '../../../../../../packages/server-utils/src/ai/core/gen-ai-attributes';
import { cleanupChildProcesses, createEsmAndCjsTests } from '../../../../utils/runner';
import { createEsmTests } from '../../../../utils/runner/createEsmAndCjsTests';
Expand Down Expand Up @@ -233,6 +235,88 @@ describe('LangChain integration (v1)', () => {
},
);

createEsmAndCjsTests(
__dirname,
'scenario-typesafe-classifier.mjs',
'instrument-with-pii.mjs',
(createRunner, test) => {
test('records a TypeSafeClassifier run as a gen_ai.evaluate span', async () => {
const runner = createRunner().ignore('event');
const spansPromise = runner.collectStreamedSpansUntilSegment('main');

await runner.start().completed();

const spans = (await spansPromise).filter(
span => span.attributes[SENTRY_ORIGIN]?.value === 'auto.ai.langchain',
);
expect(spans.map(span => span.name)).toEqual(['evaluate jev-latest', 'evaluate jev-latest']);

const [evaluateSpan, messageEvaluateSpan] = spans.sort((a, b) => a.start_timestamp - b.start_timestamp);
expect(evaluateSpan.attributes[SENTRY_OP].value).toBe(GEN_AI_EVALUATE);
expect(evaluateSpan.attributes[SENTRY_ORIGIN].value).toBe('auto.ai.langchain');
expect(evaluateSpan.attributes[GEN_AI_OPERATION_NAME].value).toBe('evaluate');
expect(evaluateSpan.attributes[GEN_AI_PROVIDER_NAME].value).toBe('typesafe');
expect(evaluateSpan.attributes[GEN_AI_REQUEST_MODEL].value).toBe('jev-latest');
expect(evaluateSpan.attributes[GEN_AI_RESPONSE_MODEL].value).toBe('jev-1.13');
expect(evaluateSpan.attributes[GEN_AI_USAGE_INPUT_TOKENS].value).toBe(30);
expect(evaluateSpan.attributes[GEN_AI_USAGE_OUTPUT_TOKENS].value).toBe(2);
expect(evaluateSpan.attributes[GEN_AI_USAGE_TOTAL_TOKENS].value).toBe(32);
expect(JSON.parse(evaluateSpan.attributes[GEN_AI_INPUT_MESSAGES].value)).toEqual([
{
type: 'evaluation',
state: 'My payouts have been failing.',
questions: { urgent: { type: 'noul', instructions: 'Is this urgent?' } },
},
]);
// A message is recorded as the transcript line the classifier sends, not as LangChain's serialized form.
expect(JSON.parse(messageEvaluateSpan!.attributes[GEN_AI_INPUT_MESSAGES].value)).toEqual([
{
type: 'evaluation',
state: 'user: My card was charged twice.',
questions: { urgent: { type: 'noul', instructions: 'Is this urgent?' } },
},
]);
expect(JSON.parse(evaluateSpan.attributes[GEN_AI_OUTPUT_MESSAGES].value)).toEqual([
{ type: 'evaluation', answers: { urgent: { type: 'noul', noul: 0.9 } } },
]);
});
},
{
additionalDependencies: {
langchain: '^1.0.0',
'@langchain/core': '^1.0.0',
'@langchain/typesafe': '^0.0.2',
},
},
);

createEsmAndCjsTests(
__dirname,
'scenario-typesafe-classifier-inherited-callbacks.mjs',
'instrument-with-pii.mjs',
(createRunner, test) => {
test('keeps the parent run callbacks for a TypeSafeClassifier called without config', async () => {
const runner = createRunner().ignore('event');
const spansPromise = runner.collectStreamedSpansUntilSegment('main');

await runner.start().completed();

const spans = await spansPromise;
const segment = spans.find(span => span.is_segment && span.name === 'main');
expect(segment!.attributes['test.recorded_runs'].value).toBe('RunnableLambda:root,TypeSafeClassifier:child');
expect(spans.filter(span => span.attributes[SENTRY_OP]?.value === GEN_AI_EVALUATE)).toHaveLength(1);
});
},
{
additionalDependencies: {
langchain: '^1.0.0',
'@langchain/core': '^1.0.0',
'@langchain/typesafe': '^0.0.2',
'@langchain/langgraph': '^1.0.0',
},
},
);

createEsmTests(
__dirname,
'scenario-openai-before-langchain.mjs',
Expand Down
26 changes: 21 additions & 5 deletions packages/server-utils/src/ai/langchain/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,11 @@ import {
import { GEN_AI_CHAT, GEN_AI_EXECUTE_TOOL, GEN_AI_INVOKE_AGENT } from '@sentry/conventions/op';
import { resolveAIRecordingOptions } from '../core/utils';
import { LANGCHAIN_ORIGIN } from './constants';
import {
addTypeSafeClassifierResponseAttributes,
isTypeSafeClassifier,
startTypeSafeClassifierSpan,
} from './typesafe-classifier';
import type {
LangChainCallbackHandler,
LangChainLLMResult,
Expand Down Expand Up @@ -43,6 +48,7 @@ export function createLangChainCallbackHandler(options: LangChainOptions = {}):

// Internal state - single instance tracks all spans
const spanMap = new Map<string, Span>();
const evaluateRunIds = new Set<string>();

/**
* Exit a span and clean up
Expand All @@ -51,8 +57,9 @@ export function createLangChainCallbackHandler(options: LangChainOptions = {}):
const span = spanMap.get(runId);
if (span?.isRecording()) {
span.end();
spanMap.delete(runId);
}
spanMap.delete(runId);
evaluateRunIds.delete(runId);
};

/**
Expand Down Expand Up @@ -215,6 +222,13 @@ export function createLangChainCallbackHandler(options: LangChainOptions = {}):
_runType?: string,
runName?: string,
) {
// A Jev call is a real model call, so it is recorded inside an agent too.
if (isTypeSafeClassifier(chain)) {
spanMap.set(runId, startTypeSafeClassifierSpan(chain, inputs, metadata, recordInputs));
evaluateRunIds.add(runId);
return;
}

// Skip chain spans when inside an agent context (createReactAgent).
// The agent already creates an invoke_agent span; internal chain steps
// (ChannelWrite, Branch, prompt, etc.) are noise.
Expand Down Expand Up @@ -264,14 +278,16 @@ export function createLangChainCallbackHandler(options: LangChainOptions = {}):
handleChainEnd(outputs: unknown, runId: string) {
const span = spanMap.get(runId);
if (span?.isRecording()) {
// Add outputs if recordOutputs is enabled
if (recordOutputs) {
if (evaluateRunIds.has(runId)) {
addTypeSafeClassifierResponseAttributes(span, outputs, recordOutputs);
} else if (recordOutputs) {
Comment thread
sentry[bot] marked this conversation as resolved.
span.setAttributes({
'langchain.chain.outputs': JSON.stringify(outputs),
});
Comment thread
sentry[bot] marked this conversation as resolved.
}
exitSpan(runId);
}
// Also for a sampled-out span, which would otherwise stay tracked.
exitSpan(runId);
},

// Chain Error Handler
Expand All @@ -281,8 +297,8 @@ export function createLangChainCallbackHandler(options: LangChainOptions = {}):
const span = spanMap.get(runId);
if (span?.isRecording()) {
span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' });
exitSpan(runId);
}
exitSpan(runId);
},

// Tool Start Handler
Expand Down
86 changes: 86 additions & 0 deletions packages/server-utils/src/ai/langchain/typesafe-classifier.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
import { SENTRY_ORIGIN } from '@sentry/conventions/attributes';
import type { Span } from '@sentry/core';
import { isObjectLike } from '@sentry/core';
import { addResponseAttributes, startEvaluateSpan } from '../typesafe';
import { LANGCHAIN_ORIGIN } from './constants';
import { getAgentNameFromMetadata, getConversationIdFromMetadata } from './utils';

// `TypeSafeClassifier` from `@langchain/typesafe` calls Jev with `fetch`, not through `@typesafe-ai/sdk`,
// so the TypeSafe integration does not see it. Its serialized id is `[...lc_namespace, lc_name()]`; only
// the class name is matched, as the namespace is the part that changes between versions.
const TYPESAFE_CLASSIFIER_NAME = 'TypeSafeClassifier';

/** The package's default, used when the classifier is constructed without a `model`. */
const DEFAULT_TYPESAFE_CLASSIFIER_MODEL = 'jev-latest';

// The `state` the classifier sends to Jev (messages rendered as transcript lines), keyed by the input
// passed to `invoke()`. LangChain hands that same object to `handleChainStart`, so the span records what
// Jev received rather than LangChain's serialized messages.
const wireStates = new WeakMap<object, unknown>();

/** Record the `state` a `TypeSafeClassifier` sends for `input`, using the classifier's own serializer. */
export function recordTypeSafeClassifierState(classifier: unknown, input: unknown): void {
// A string state is sent as is, so there is nothing to record.
if (!isObjectLike(input) || !isObjectLike(classifier) || typeof classifier.payload !== 'function') {
return;
}

try {
const body: unknown = JSON.parse(classifier.payload(input));
if (isObjectLike(body)) {
wireStates.set(input, body.state);
}
} catch {
// The classifier rejects the same input itself; the span keeps LangChain's form of it.
}
}

// Typed loosely: LangChain's `Serialized` union does not match our handler's chain type.
export function isTypeSafeClassifier(chain: unknown): boolean {
return isObjectLike(chain) && Array.isArray(chain.id) && chain.id[chain.id.length - 1] === TYPESAFE_CLASSIFIER_NAME;
}

/** Start an `evaluate` span for a `TypeSafeClassifier` run, from its serialized constructor arguments. */
export function startTypeSafeClassifierSpan(
chain: unknown,
inputs: Record<string, unknown>,
metadata: Record<string, unknown> | undefined,
recordInputs: boolean,
): Span {
const kwargs = isObjectLike(chain) && isObjectLike(chain.kwargs) ? chain.kwargs : {};
const model = typeof kwargs.model === 'string' ? kwargs.model : DEFAULT_TYPESAFE_CLASSIFIER_MODEL;

return startEvaluateSpan({ model, state: getState(inputs), questions: kwargs.questions }, undefined, recordInputs, {
[SENTRY_ORIGIN]: LANGCHAIN_ORIGIN,
...getAgentNameFromMetadata(metadata),
...getConversationIdFromMetadata(metadata),
});
}

/** The classifier result reports usage in camelCase, the TypeSafe helpers read the API's snake_case. */
export function addTypeSafeClassifierResponseAttributes(span: Span, outputs: unknown, recordOutputs: boolean): void {
if (!isObjectLike(outputs)) {
return;
}

const usage = isObjectLike(outputs.usage) ? outputs.usage : {};
addResponseAttributes(
span,
{
model: outputs.model,
answers: outputs.answers,
usage: { input_tokens: usage.inputTokens, output_tokens: usage.outputTokens },
},
recordOutputs,
);
}

function getState(inputs: Record<string, unknown>): unknown {
// LangChain hands a string or array input to callbacks wrapped as `{ input }`.
const keys = Object.keys(inputs);
const input = inputs.input;
const state =
keys.length === 1 && keys[0] === 'input' && (typeof input === 'string' || Array.isArray(input)) ? input : inputs;

return isObjectLike(state) && wireStates.has(state) ? wireStates.get(state) : state;
}
16 changes: 16 additions & 0 deletions packages/server-utils/src/ai/langchain/utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -564,6 +564,22 @@ export function extractToolDefinitions(extraParams?: Record<string, unknown>): s
return JSON.stringify(toolDefs);
}

// A run started without config inherits the parent run's config, which LangChain keeps in a global
// AsyncLocalStorage. `ensureConfig` lets an explicit `callbacks` key replace the inherited one.
const LANGCHAIN_ALS_KEY = Symbol.for('ls:tracing_async_local_storage');
const LANGCHAIN_CHILD_CONFIG_KEY = Symbol.for('lc:child_config');

/** The callbacks LangChain gives a run started without config: the ones of the parent run, if any. */
export function getInheritedLangChainCallbacks(): unknown {
const storage = (globalThis as Record<symbol, unknown>)[LANGCHAIN_ALS_KEY] as
| { getStore?: () => unknown }
| undefined;
const store = typeof storage?.getStore === 'function' ? storage.getStore() : undefined;
const extra = isObjectLike(store) ? (store.extra as Record<symbol, unknown> | undefined) : undefined;
const config = isObjectLike(extra) ? extra[LANGCHAIN_CHILD_CONFIG_KEY] : undefined;
return isObjectLike(config) ? config.callbacks : undefined;
}

/** Duck-types a LangChain `CallbackManager` (avoids coupling to a specific `@langchain/core` resolution). */
function isCallbackManager(value: unknown): value is {
addHandler: (handler: unknown, inherit?: boolean) => void;
Expand Down
Loading
Loading