Skip to content
Open
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
38 changes: 37 additions & 1 deletion src/langfuse.ts
Original file line number Diff line number Diff line change
Expand Up @@ -981,10 +981,12 @@ export class LangfuseClient {

this.ensureGenerationParent(input.sessionID);

const observationName = this.getToolObservationName(input.tool, input.args);

this.withObservationParent(
input.sessionID,
() => {
const span = this.traceState.tracer.startSpan(input.tool, {
const span = this.traceState.tracer.startSpan(observationName, {
attributes: {
"langfuse.observation.type": "tool",
"session.id": input.sessionID,
Expand Down Expand Up @@ -1129,6 +1131,40 @@ export class LangfuseClient {
this.traceState.toolMessageIdsByCallId.delete(input.callID);
}

// Use skill or subagent names when available; otherwise fall back to the tool name.
private getToolObservationName(tool: string, args: unknown) {
if (typeof args !== "object" || args === null || Array.isArray(args)) {
return tool;
}

let semanticName: unknown;
if (tool === "skill") {
// OpenCode v1 sends the skill name; OpenCode v2 sends the skill id.
if ("name" in args && typeof args.name === "string") {
semanticName = args.name;
} else if ("id" in args) {
semanticName = args.id;
} else {
return tool;
}
// OpenCode v1 represents subagent calls as task tools.
} else if (tool === "task" && "subagent_type" in args) {
semanticName = args.subagent_type;
// OpenCode v2 represents subagent calls as dedicated subagent tools.
} else if (tool === "subagent" && "agent" in args) {
semanticName = args.agent;
} else {
return tool;
}

if (typeof semanticName !== "string") {
return tool;
}

const name = semanticName.trim();
return name === "" ? tool : `${tool}:${name}`;
}

private ensureGenerationParent(sessionID: string) {
if (
this.traceState.activeGenerationSteps.has(sessionID) ||
Expand Down
226 changes: 226 additions & 0 deletions test/integration/v1.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1108,6 +1108,232 @@ describe("built plugin", { concurrent: false }, () => {
).toHaveLength(0);
});

test.each([
{
tool: "skill",
args: { name: "resolve-dependencies" },
observationName: "skill:resolve-dependencies",
},
{
tool: "task",
args: { subagent_type: "developer", prompt: "Fix the build" },
observationName: "task:developer",
},
])(
"exports $observationName without changing tool data",
async ({ tool, args, observationName }) => {
const sessionID = "semantic-tool-session";
const callID = "semantic-tool-call";
const messageID = "semantic-tool-assistant";
await sendUserMessage({
sessionID,
messageID: "semantic-tool-user",
text: "Fix the build",
started: startedAt,
});
await startGeneration({
id: "semantic-tool-step",
sessionID,
assistantMessageID: messageID,
started: startedAt + 100,
});
await hooks["tool.execute.before"]?.(
{ sessionID, callID, tool },
{ args },
);
const part = {
id: "semantic-tool-part",
sessionID,
messageID,
type: "tool" as const,
callID,
tool,
};
await emitEvent({
type: "message.part.updated",
properties: {
part: {
...part,
state: {
status: "running",
input: args,
time: { start: startedAt + 200 },
},
},
},
});
await hooks["tool.execute.after"]?.(
{ sessionID, callID, tool, args },
{ title: "Done", output: "ok", metadata: {} },
);
await emitEvent({
type: "message.part.updated",
properties: {
part: {
...part,
state: {
status: "completed",
input: args,
title: "Done",
output: "ok",
metadata: {},
time: { start: startedAt + 200, end: startedAt + 300 },
},
},
},
});

const { spans } = await flushSession(sessionID);
const tools = spans.filter(
(span) => getAttributes(span)["langfuse.observation.type"] === "tool",
);
expect(tools).toHaveLength(1);
const observation = tools[0];
expect(observation.name).toBe(observationName);
expect(observation.parentSpanId).toBe(
getSpan(spans, "opencode.generation").spanId,
);
expect(
getJsonAttribute(observation, "langfuse.observation.input"),
).toEqual(args);
expect(
getJsonAttribute(observation, "langfuse.observation.metadata"),
).toEqual({ callID, tool });
expect(
getJsonAttribute(observation, "langfuse.observation.output"),
).toEqual({ title: "Done", output: "ok" });
},
);

test.each(["completed", "error"] as const)(
"names semantic tools from %s parts without execution hooks",
async (status) => {
const sessionID = "semantic-tool-parts-session";
for (const [tool, args, observationName] of [
[
"skill",
{ name: "resolve-dependencies" },
"skill:resolve-dependencies",
],
["task", { subagent_type: "developer" }, "task:developer"],
] as const) {
await emitEvent({
type: "message.part.updated",
properties: {
part: {
id: `${tool}-part`,
sessionID,
messageID: "semantic-tool-parts-assistant",
type: "tool",
callID: `${tool}-call`,
tool,
state: {
input: args,
time: { start: startedAt, end: startedAt + 100 },
...(status === "completed"
? { status, title: "Done", output: "ok", metadata: {} }
: { status, error: "Tool failed" }),
},
},
},
});
const { spans } = await flushSession(sessionID);
expect(spans).toHaveLength(1);
const observation = getSpan(spans, observationName);
expect(
getJsonAttribute(observation, "langfuse.observation.input"),
).toEqual(args);
expect(
getJsonAttribute(observation, "langfuse.observation.metadata"),
).toEqual({ callID: `${tool}-call`, tool });
if (status === "error") {
expect(observation.status?.code).toBe(2);
expect(
getJsonAttribute(observation, "langfuse.observation.output"),
).toEqual({ error: "Tool failed" });
}
}
},
);

test("keeps fallback tool names and preserves inputs when normalizing names", async () => {
const sessionID = "tool-name-fallback-session";
const cases = [
{ tool: "skill", args: {}, name: "skill" },
{ tool: "skill", args: { name: "" }, name: "skill" },
{ tool: "skill", args: { name: " \t " }, name: "skill" },
{ tool: "skill", args: { name: 42 }, name: "skill" },
{ tool: "skill", args: { name: ["review"] }, name: "skill" },
{ tool: "skill", args: { id: "" }, name: "skill" },
{ tool: "skill", args: { id: 42 }, name: "skill" },
{ tool: "skill", args: null, name: "skill" },
{ tool: "skill", args: "review", name: "skill" },
{ tool: "task", args: { name: "developer" }, name: "task" },
{ tool: "task", args: { subagent_type: "" }, name: "task" },
{ tool: "task", args: { subagent_type: " \n " }, name: "task" },
{ tool: "task", args: { subagent_type: null }, name: "task" },
{ tool: "task", args: [], name: "task" },
{ tool: "subagent", args: {}, name: "subagent" },
{ tool: "subagent", args: { agent: "" }, name: "subagent" },
{ tool: "subagent", args: { agent: 42 }, name: "subagent" },
{
tool: "read",
args: { name: "README.md", subagent_type: "developer" },
name: "read",
},
{ tool: "skill", args: { name: " review " }, name: "skill:review" },
{
tool: "skill",
args: { id: " build-project " },
name: "skill:build-project",
},
{
tool: "task",
args: { subagent_type: " developer " },
name: "task:developer",
},
{
tool: "subagent",
args: { agent: " ts-reviewer " },
name: "subagent:ts-reviewer",
},
];
for (const [index, { tool, args }] of cases.entries()) {
await hooks["tool.execute.after"]?.(
{ sessionID, callID: `fallback-${index.toString()}`, tool, args },
{ title: "Done", output: "ok", metadata: {} },
);
}
const { spans } = await flushSession(sessionID);
expect(spans).toHaveLength(cases.length);
for (const [index, { tool, args, name }] of cases.entries()) {
const callID = `fallback-${index.toString()}`;
const observation = spans.find((span) => {
const metadata = getJsonAttribute(
span,
"langfuse.observation.metadata",
);
return (
typeof metadata === "object" &&
metadata !== null &&
"callID" in metadata &&
metadata.callID === callID
);
});
expect(observation).toBeDefined();
if (!observation) {
throw new Error(`Expected tool observation ${callID}`);
}
expect(observation.name).toBe(name);
expect(
getJsonAttribute(observation, "langfuse.observation.input"),
).toEqual(args);
expect(
getJsonAttribute(observation, "langfuse.observation.metadata"),
).toEqual({ callID, tool });
}
});

// https://github.com/anomalyco/opencode/blob/v1.15.13/packages/core/src/session-event.ts#L353-L362
test("supports OpenCode >=1.15.13 <1.16 compaction events", async () => {
const sessionID = "legacy-compaction-session";
Expand Down
88 changes: 88 additions & 0 deletions test/integration/v2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -293,6 +293,94 @@ describe("OpenCode 2 package entrypoint", () => {
await cleanup?.();
});

test.each([
{
tool: "skill",
input: { id: "build-project" },
},
{
tool: "task",
input: { subagent_type: "developer", prompt: "Fix the build" },
},
{
tool: "subagent",
input: { agent: "ts-reviewer", prompt: "Review the build" },
},
])(
"forwards semantic $tool input for observation naming",
async (toolCall) => {
let executeBefore: ((input: unknown) => void) | undefined;
let executeAfter: ((input: unknown) => void) | undefined;
const registration = { dispose: vi.fn(() => Promise.resolve()) };
const contextInput: unknown = {
app: { version: "2.0.4" },
session: { hook: vi.fn(() => Promise.resolve(registration)) },
tool: {
hook: vi.fn((name: string, handler: (input: unknown) => void) => {
if (name === "execute.before") {
executeBefore = handler;
}
if (name === "execute.after") {
executeAfter = handler;
}
return Promise.resolve(registration);
}),
},
event: {
subscribe: () => ({
async *[Symbol.asyncIterator]() {
await Promise.resolve();
yield* [];
},
}),
},
};
const context = Schema.decodeUnknownSync(
Schema.declare(
(input): input is Parameters<typeof SourcePlugin.setup>[0] =>
typeof input === "object" && input !== null,
),
)(contextInput);

const cleanup = await SourcePlugin.setup(context);
expect(executeBefore).toBeTypeOf("function");
expect(executeAfter).toBeTypeOf("function");

const input = {
id: `${toolCall.tool}-call`,
messageID: "assistant-1",
sessionID: "session-1",
tool: toolCall.tool,
input: toolCall.input,
};
executeBefore?.(input);
executeAfter?.({
...input,
status: "success",
result: { content: "ok" },
});

expect(runtime.traceToolStart).toHaveBeenCalledWith({
sessionID: "session-1",
messageID: "assistant-1",
callID: `${toolCall.tool}-call`,
tool: toolCall.tool,
args: toolCall.input,
});
expect(runtime.traceToolEnd).toHaveBeenCalledWith({
sessionID: "session-1",
messageID: "assistant-1",
callID: `${toolCall.tool}-call`,
tool: toolCall.tool,
args: toolCall.input,
title: toolCall.tool,
output: "ok",
});

await cleanup?.();
},
);

test("captures the complete model input from each context hook", async () => {
let context:
| ((input: {
Expand Down