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
92 changes: 69 additions & 23 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,40 +39,74 @@ const EMIT_IMAGE_MEDIA = (() => {
return raw ? !["false", "0"].includes(raw) : true;
})();

const ENV_TRACEPARENT = "LANGFUSE_PI_TRACEPARENT";
const ENV_PARENT_TRACE_ID = "LANGFUSE_PI_PARENT_TRACE_ID";
const ENV_PARENT_SPAN_ID = "LANGFUSE_PI_PARENT_SPAN_ID";
const ENV_PARENT_SESSION_ID = "LANGFUSE_PI_PARENT_SESSION_ID";
const ENV_PARENT_DEPTH = "LANGFUSE_PI_PARENT_DEPTH";
const ENV_PARENT_EXTERNAL_TRACE = "LANGFUSE_PI_PARENT_EXTERNAL_TRACE";

const HEX_TRACE_ID = /^[0-9a-f]{32}$/;
const HEX_SPAN_ID = /^[0-9a-f]{16}$/;

const isUsableTraceId = (v: string | undefined): v is string =>
!!v && HEX_TRACE_ID.test(v) && !/^0+$/.test(v);
const isUsableSpanId = (v: string | undefined): v is string =>
!!v && HEX_SPAN_ID.test(v) && !/^0+$/.test(v);

const TRACEPARENT = /^(?!ff)([0-9a-f]{2})-([0-9a-f]{32})-([0-9a-f]{16})-([0-9a-f]{2})(-.*)?$/;

export function parseTraceparent(value: string): { traceId: string; spanId: string } | undefined {
const match = TRACEPARENT.exec(value.trim().toLowerCase());
if (!match) return undefined;
const [, version, traceId, spanId, , trailing] = match;
if (version === "00" && trailing) return undefined;
if (!isUsableTraceId(traceId) || !isUsableSpanId(spanId)) return undefined;
return { traceId, spanId };
}

export interface InheritedParent {
spanContext: SpanContext;
sessionId?: string;
depth: number;
source: "subagent" | "attached";
externalTrace: boolean;
}

export function readInheritedParent(env: NodeJS.ProcessEnv = process.env): InheritedParent | undefined {
const traceId = env[ENV_PARENT_TRACE_ID]?.trim().toLowerCase();
const spanId = env[ENV_PARENT_SPAN_ID]?.trim().toLowerCase();
if (!traceId || !spanId) return undefined;
if (!HEX_TRACE_ID.test(traceId) || !HEX_SPAN_ID.test(spanId)) return undefined;
// OTel defines the all-zero ids as invalid.
if (/^0+$/.test(traceId) || /^0+$/.test(spanId)) return undefined;
const parent = readParentIds(env);
if (!parent) return undefined;
const depth = Number(env[ENV_PARENT_DEPTH] ?? "0");
return {
spanContext: {
traceId,
spanId,
traceId: parent.traceId,
spanId: parent.spanId,
traceFlags: TraceFlags.SAMPLED,
isRemote: true,
},
sessionId: env[ENV_PARENT_SESSION_ID]?.trim() || undefined,
depth: Number.isFinite(depth) && depth > 0 ? depth : 0,
source: parent.source,
externalTrace: parent.source === "attached" || env[ENV_PARENT_EXTERNAL_TRACE] === "1",
};
}

type ParentIds = { traceId: string; spanId: string; source: "subagent" | "attached" };

function readParentIds(env: NodeJS.ProcessEnv): ParentIds | undefined {
const raw = env[ENV_TRACEPARENT]?.trim();
if (raw) {
const parsed = parseTraceparent(raw);
if (parsed) return { ...parsed, source: "attached" };
console.error(`[pi-langfuse] Ignoring malformed ${ENV_TRACEPARENT}: ${JSON.stringify(raw)}`);
}
const traceId = env[ENV_PARENT_TRACE_ID]?.trim().toLowerCase();
const spanId = env[ENV_PARENT_SPAN_ID]?.trim().toLowerCase();
if (!traceId && !spanId) return undefined;
if (!isUsableTraceId(traceId) || !isUsableSpanId(spanId)) return undefined;
return { traceId, spanId, source: "subagent" };
}

// ---------------------------------------------------------------------------
// Config
// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -661,11 +695,15 @@ export default function (pi: ExtensionAPI) {
let lastContextHistory: ChatMlMessage[] | undefined;
let compactionStartedAt: Date | undefined;
const inheritedParent = readInheritedParent();
const isAttached = inheritedParent?.source === "attached";
const externalTrace = inheritedParent?.externalTrace ?? false;
const inheritedParentEnv: Record<string, string | undefined> = {
[ENV_TRACEPARENT]: process.env[ENV_TRACEPARENT],
[ENV_PARENT_TRACE_ID]: process.env[ENV_PARENT_TRACE_ID],
[ENV_PARENT_SPAN_ID]: process.env[ENV_PARENT_SPAN_ID],
[ENV_PARENT_SESSION_ID]: process.env[ENV_PARENT_SESSION_ID],
[ENV_PARENT_DEPTH]: process.env[ENV_PARENT_DEPTH],
[ENV_PARENT_EXTERNAL_TRACE]: process.env[ENV_PARENT_EXTERNAL_TRACE],
};

const ensureRuntime = (): Runtime => {
Expand Down Expand Up @@ -718,10 +756,13 @@ export default function (pi: ExtensionAPI) {
// Do not publish a root that the sampler dropped. Child spans would point
// to a trace with no exported root.
if (!(ctx.traceFlags & TraceFlags.SAMPLED)) return;
delete process.env[ENV_TRACEPARENT];
process.env[ENV_PARENT_TRACE_ID] = ctx.traceId;
process.env[ENV_PARENT_SPAN_ID] = ctx.spanId;
process.env[ENV_PARENT_SESSION_ID] = sessionId;
process.env[ENV_PARENT_DEPTH] = String((inheritedParent?.depth ?? 0) + 1);
if (externalTrace) process.env[ENV_PARENT_EXTERNAL_TRACE] = "1";
else delete process.env[ENV_PARENT_EXTERNAL_TRACE];
};

// A child must not attach to a turn that has ended. The inherited values
Expand Down Expand Up @@ -805,15 +846,17 @@ export default function (pi: ExtensionAPI) {
const userText = [event.prompt, ...promptImages.map(describeImage)].filter(Boolean).join("\n");
lastPromptText = userText;
lastContextHistory = undefined;
const isSubagent = Boolean(inheritedParent);
traceAttributes = {
...(isSubagent ? {} : { [LangfuseOtelSpanAttributes.TRACE_NAME]: TRACE_NAME }),
[LangfuseOtelSpanAttributes.TRACE_SESSION_ID]: isSubagent
? (inheritedParent!.sessionId ?? sessionId)
: sessionId,
[LangfuseOtelSpanAttributes.TRACE_TAGS]: BASE_TAGS,
...(config.userId ? { [LangfuseOtelSpanAttributes.TRACE_USER_ID]: config.userId } : {}),
};
const isSubagent = inheritedParent?.source === "subagent";
traceAttributes = externalTrace
? {}
: {
...(isSubagent ? {} : { [LangfuseOtelSpanAttributes.TRACE_NAME]: TRACE_NAME }),
[LangfuseOtelSpanAttributes.TRACE_SESSION_ID]: isSubagent
? (inheritedParent!.sessionId ?? sessionId)
: sessionId,
[LangfuseOtelSpanAttributes.TRACE_TAGS]: BASE_TAGS,
...(config.userId ? { [LangfuseOtelSpanAttributes.TRACE_USER_ID]: config.userId } : {}),
};

const root = startObservation(
isSubagent ? SUBAGENT_ROOT_OBSERVATION_NAME : ROOT_OBSERVATION_NAME,
Expand All @@ -831,6 +874,7 @@ export default function (pi: ExtensionAPI) {
...(isSubagent
? { pi_subagent: true, subagent_depth: inheritedParent!.depth, parent_session_id: inheritedParent!.sessionId }
: {}),
...(isAttached ? { attached_to_external_parent: true } : {}),
},
},
{ asType: "span", ...(inheritedParent ? { parentSpanContext: inheritedParent.spanContext } : {}) },
Expand Down Expand Up @@ -1081,12 +1125,14 @@ export default function (pi: ExtensionAPI) {
// /compact in a fresh process would hit the unset tracer provider.
ensureRuntime();
const sessionId = ctx.sessionManager.getSessionId();
traceAttributes = {
...(inheritedParent ? {} : { [LangfuseOtelSpanAttributes.TRACE_NAME]: `Pi ${name}` }),
[LangfuseOtelSpanAttributes.TRACE_SESSION_ID]: inheritedParent?.sessionId ?? sessionId,
[LangfuseOtelSpanAttributes.TRACE_TAGS]: BASE_TAGS,
...(config.userId ? { [LangfuseOtelSpanAttributes.TRACE_USER_ID]: config.userId } : {}),
};
traceAttributes = externalTrace
? {}
: {
...(inheritedParent ? {} : { [LangfuseOtelSpanAttributes.TRACE_NAME]: `Pi ${name}` }),
[LangfuseOtelSpanAttributes.TRACE_SESSION_ID]: inheritedParent?.sessionId ?? sessionId,
[LangfuseOtelSpanAttributes.TRACE_TAGS]: BASE_TAGS,
...(config.userId ? { [LangfuseOtelSpanAttributes.TRACE_USER_ID]: config.userId } : {}),
};
const obs = startObservation(
name,
{
Expand Down
3 changes: 3 additions & 0 deletions test/fixtures/env-probe.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,11 @@ import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
export default function (pi: ExtensionAPI) {
pi.on("before_provider_request", async () => {
console.error(`PROBE turn ${process.env.LANGFUSE_PI_PARENT_TRACE_ID ?? "<unset>"}`);
console.error(`PROBE turn-tp ${process.env.LANGFUSE_PI_TRACEPARENT ?? "<unset>"}`);
console.error(`PROBE turn-ext ${process.env.LANGFUSE_PI_PARENT_EXTERNAL_TRACE ?? "<unset>"}`);
});
pi.on("agent_settled", async () => {
console.error(`PROBE settled ${process.env.LANGFUSE_PI_PARENT_TRACE_ID ?? "<unset>"}`);
console.error(`PROBE settled-tp ${process.env.LANGFUSE_PI_TRACEPARENT ?? "<unset>"}`);
});
}
9 changes: 2 additions & 7 deletions test/helpers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -223,7 +223,6 @@ export interface Capture {
port: number;
requests: unknown[];
spans: () => CapturedSpan[];
/** Attributes of the OTLP resource the exported spans were sent under. */
resourceAttrs: () => Record<string, unknown>;
close: () => void;
}
Expand Down Expand Up @@ -318,7 +317,6 @@ export interface Sandbox {
export interface SandboxOptions {
contextWindow?: number;
keepRecentTokens?: number;
/** Pads README.md so reading it grows the context past keepRecentTokens. */
readmeFillerLines?: number;
}

Expand Down Expand Up @@ -378,7 +376,6 @@ export function runPi(
opts: {
continue?: boolean;
env?: Record<string, string | undefined>;
/** More extensions to load with ours, for example the subagent fixture. */
extensions?: string[];
} = {},
): Promise<{ status: number | null; stdout: string; stderr: string }> {
Expand All @@ -399,7 +396,6 @@ export function runPi(
return new Promise((resolvePromise) => {
const child = spawn(PI_BIN, args, {
cwd: sandbox.workspace,
// stdin must be closed — pi's print mode waits for EOF on piped stdin.
stdio: ["ignore", "pipe", "pipe"],
env: {
...process.env,
Expand All @@ -414,15 +410,15 @@ export function runPi(
LANGFUSE_TRACING_ENVIRONMENT: undefined,
LANGFUSE_RELEASE: undefined,
LANGFUSE_TRACING_ENABLED: undefined,
// These feed the exported OTel resource, so the runner's own
// environment must not reach the span payload either.
OTEL_SERVICE_NAME: undefined,
OTEL_RESOURCE_ATTRIBUTES: undefined,
// A test run must not get a parent trace from this test process.
LANGFUSE_PI_TRACEPARENT: undefined,
LANGFUSE_PI_PARENT_TRACE_ID: undefined,
LANGFUSE_PI_PARENT_SPAN_ID: undefined,
LANGFUSE_PI_PARENT_SESSION_ID: undefined,
LANGFUSE_PI_PARENT_DEPTH: undefined,
LANGFUSE_PI_PARENT_EXTERNAL_TRACE: undefined,
// The subagent fixture uses these values.
TEST_PI_BIN: PI_BIN,
TEST_LANGFUSE_EXTENSION: EXTENSION,
Expand All @@ -448,7 +444,6 @@ export function runPi(
});
}

/** Poll until the capture holds at least `n` export requests (flushes are async). */
export async function waitForRequests(capture: Capture, n: number, timeoutMs = 5000): Promise<void> {
const start = Date.now();
while (capture.requests.length < n && Date.now() - start < timeoutMs) {
Expand Down
Loading
Loading