Compare commits

...

3 Commits

Author SHA1 Message Date
abeatrix ce802b2ba7 use observeOpenAI 2025-12-15 19:46:07 -08:00
abeatrix 8515611dff Fix Otel versions 2025-12-15 17:31:36 -08:00
abeatrix 33793a6c7e poc feat: add Langfuse integration for LLM tracing
- Add Langfuse client, otel, and tracing packages as dependencies
- Update OpenTelemetry packages to version 0.202.0 for compatibility
- Add Langfuse configuration options to .env.example
- Enable advanced LLM observability through Langfuse integration

Note: This  is a POC only and only basic llm call tracing info for each task is implemented. The span will not be available until the task is completed.
2025-12-15 17:27:23 -08:00
11 changed files with 742 additions and 832 deletions
+7
View File
@@ -91,6 +91,13 @@ POSTHOG_TELEMETRY_ENABLED=true # Enable PostHog telemetry (default: tru
# OTEL_EXPORTER_OTLP_ENDPOINT=https://otel.example.com
# OTEL_EXPORTER_OTLP_HEADERS=authorization=Bearer your-token
# ============================================================================
# Langfuse Integration (Optional - for advanced LLM tracing)
# ============================================================================
# LANGFUSE_SECRET_KEY = "sk-lf-..."
# LANGFUSE_PUBLIC_KEY = "pk-lf-..."
# LANGFUSE_BASE_URL = "https://cloud.langfuse.com" # Default to "https://us.cloud.langfuse.com"
# ============================================================================
# OPTIONAL DEVELOPMENT SETTINGS
# ============================================================================
+626 -799
View File
File diff suppressed because it is too large Load Diff
+17 -13
View File
@@ -462,26 +462,30 @@
"@google/genai": "^1.30.0",
"@grpc/grpc-js": "^1.9.15",
"@grpc/reflection": "^1.0.4",
"@langfuse/client": "^4.4.9",
"@langfuse/openai": "^4.4.9",
"@langfuse/otel": "^4.4.9",
"@langfuse/tracing": "^4.4.9",
"@mistralai/mistralai": "^1.5.0",
"@modelcontextprotocol/sdk": "^1.11.1",
"@opentelemetry/api": "^1.9.0",
"@opentelemetry/core": "^2.1.0",
"@opentelemetry/exporter-logs-otlp-grpc": "^0.56.0",
"@opentelemetry/exporter-logs-otlp-http": "^0.56.0",
"@opentelemetry/exporter-logs-otlp-proto": "^0.56.0",
"@opentelemetry/exporter-metrics-otlp-grpc": "^0.56.0",
"@opentelemetry/exporter-metrics-otlp-http": "^0.56.0",
"@opentelemetry/exporter-metrics-otlp-proto": "^0.56.0",
"@opentelemetry/exporter-prometheus": "^0.56.0",
"@opentelemetry/exporter-trace-otlp-http": "^0.56.0",
"@opentelemetry/exporter-logs-otlp-grpc": "^0.202.0",
"@opentelemetry/exporter-logs-otlp-http": "^0.202.0",
"@opentelemetry/exporter-logs-otlp-proto": "^0.202.0",
"@opentelemetry/exporter-metrics-otlp-grpc": "^0.202.0",
"@opentelemetry/exporter-metrics-otlp-http": "^0.202.0",
"@opentelemetry/exporter-metrics-otlp-proto": "^0.202.0",
"@opentelemetry/exporter-prometheus": "^0.202.0",
"@opentelemetry/exporter-trace-otlp-http": "^0.202.0",
"@opentelemetry/instrumentation": "^0.205.0",
"@opentelemetry/instrumentation-http": "^0.205.0",
"@opentelemetry/resources": "^1.30.1",
"@opentelemetry/sdk-logs": "^0.56.0",
"@opentelemetry/sdk-metrics": "^1.30.1",
"@opentelemetry/sdk-node": "^0.56.0",
"@opentelemetry/resources": "^2.1.0",
"@opentelemetry/sdk-logs": "^0.202.0",
"@opentelemetry/sdk-metrics": "^2.1.0",
"@opentelemetry/sdk-node": "^0.202.0",
"@opentelemetry/sdk-trace-base": "^2.1.0",
"@opentelemetry/sdk-trace-node": "^1.30.1",
"@opentelemetry/sdk-trace-node": "^2.1.0",
"@opentelemetry/semantic-conventions": "^1.37.0",
"@playwright/test": "^1.55.1",
"@sap-ai-sdk/ai-api": "^2.1.0",
+2 -1
View File
@@ -1,3 +1,4 @@
import { LangfuseSpan } from "@langfuse/tracing"
import { ApiConfiguration, ModelInfo, QwenApiRegions } from "@shared/api"
import { Mode } from "@shared/storage/types"
import { ClineStorageMessage } from "@/shared/messages/content"
@@ -48,7 +49,7 @@ export type CommonApiHandlerOptions = {
onRetryAttempt?: ApiConfiguration["onRetryAttempt"]
}
export interface ApiHandler {
createMessage(systemPrompt: string, messages: ClineStorageMessage[], tools?: ClineTool[], useResponseApi?: boolean): ApiStream
createMessage(systemPrompt: string, messages: ClineStorageMessage[], tools?: ClineTool[], span?: LangfuseSpan): ApiStream
getModel(): ApiHandlerModel
getApiStreamUsage?(): Promise<ApiStreamUsageChunk | undefined>
abort?(): void
+43 -11
View File
@@ -1,3 +1,5 @@
import { observeOpenAI } from "@langfuse/openai"
import { LangfuseSpan, updateActiveTrace } from "@langfuse/tracing"
import { ModelInfo, openRouterDefaultModelId, openRouterDefaultModelInfo } from "@shared/api"
import { shouldSkipReasoningForModel } from "@utils/model-utils"
import axios from "axios"
@@ -98,20 +100,38 @@ export class ClineHandler implements ApiHandler {
}
@withRetry()
async *createMessage(systemPrompt: string, messages: ClineStorageMessage[], tools?: OpenAITool[]): ApiStream {
async *createMessage(
systemPrompt: string,
messages: ClineStorageMessage[],
tools?: OpenAITool[],
span?: LangfuseSpan,
): ApiStream {
try {
const client = await this.ensureClient()
this.lastGenerationId = undefined
this.lastRequestId = undefined
let didOutputUsage: boolean = false
const openaiClient = await this.ensureClient()
const client = observeOpenAI(openaiClient, {
generationName: "cline",
sessionId: this.options.ulid,
})
const model = this.getModel()
const toolCallProcessor = new ToolCallProcessor()
const lastUserContent = messages.at(-1)?.content
const input =
typeof lastUserContent === "string"
? lastUserContent
: lastUserContent?.map((m) => (m.type === "text" ? m.text : "")).join("\n\n")
span?.update({ input, metadata: { model: model.id } })
const stream = await createOpenRouterStream(
client,
systemPrompt,
messages,
this.getModel(),
model,
this.options.reasoningEffort,
this.options.thinkingBudgetTokens,
this.options.openRouterProviderSorting,
@@ -119,8 +139,6 @@ export class ClineHandler implements ApiHandler {
this.options.geminiThinkingLevel,
)
const toolCallProcessor = new ToolCallProcessor()
for await (const chunk of stream) {
Logger.debug("ClineHandler chunk:" + JSON.stringify(chunk))
// openrouter returns an error object instead of the openai sdk throwing an error
@@ -203,12 +221,15 @@ export class ClineHandler implements ApiHandler {
totalCost = 0
}
const inputTokens = (chunk.usage.prompt_tokens || 0) - (chunk.usage.prompt_tokens_details?.cached_tokens || 0)
const outputTokens = chunk.usage.completion_tokens || 0
yield {
type: "usage",
cacheWriteTokens: 0,
cacheReadTokens: chunk.usage.prompt_tokens_details?.cached_tokens || 0,
inputTokens: (chunk.usage.prompt_tokens || 0) - (chunk.usage.prompt_tokens_details?.cached_tokens || 0),
outputTokens: chunk.usage.completion_tokens || 0,
inputTokens,
outputTokens,
totalCost,
}
didOutputUsage = true
@@ -226,6 +247,11 @@ export class ClineHandler implements ApiHandler {
} catch (error) {
console.error("Cline API Error:", error)
throw error
} finally {
updateActiveTrace({
userId: "test",
sessionId: this.options.ulid,
})
}
}
@@ -249,14 +275,20 @@ export class ClineHandler implements ApiHandler {
})
const generation = response.data
const totalCost = generation?.total_cost || 0
const inputTokens =
(generation.usage.prompt_tokens || 0) - (generation.usage.prompt_tokens_details?.cached_tokens || 0)
const outputTokens = generation.usage.completion_tokens || 0
return {
type: "usage",
cacheWriteTokens: 0,
cacheReadTokens: generation?.native_tokens_cached || 0,
// openrouter generation endpoint fails often
inputTokens: (generation?.native_tokens_prompt || 0) - (generation?.native_tokens_cached || 0),
outputTokens: generation?.native_tokens_completion || 0,
totalCost: generation?.total_cost || 0,
inputTokens,
outputTokens,
totalCost,
}
} catch (error) {
// ignore if fails
+5 -1
View File
@@ -1,3 +1,4 @@
import { observeOpenAI } from "@langfuse/openai"
import {
ModelInfo,
OpenAiCompatibleModelInfo,
@@ -86,7 +87,10 @@ export class OpenAiNativeHandler implements ApiHandler {
messages: ClineStorageMessage[],
tools?: ChatCompletionTool[],
): ApiStream {
const client = this.ensureClient()
const openaiClient = this.ensureClient()
const client = observeOpenAI(openaiClient, {
generationName: "openai-native",
})
const model = this.getModel()
const toolCallProcessor = new ToolCallProcessor()
+5 -1
View File
@@ -1,3 +1,4 @@
import { observeOpenAI } from "@langfuse/openai"
import { ModelInfo, openRouterDefaultModelId, openRouterDefaultModelInfo } from "@shared/api"
import OpenAI from "openai"
import type { ChatCompletionTool as OpenAITool } from "openai/resources/chat/completions"
@@ -48,7 +49,10 @@ export class VercelAIGatewayHandler implements ApiHandler {
@withRetry()
async *createMessage(systemPrompt: string, messages: ClineStorageMessage[], tools?: OpenAITool[]): ApiStream {
const client = this.ensureClient()
const openaiClient = this.ensureClient()
const client = observeOpenAI(openaiClient, {
generationName: "vercel-ai-gateway",
})
const modelId = this.getModel().id
const modelInfo = this.getModel().info
+14 -1
View File
@@ -44,6 +44,7 @@ import { formatContentBlockToMarkdown } from "@integrations/misc/export-markdown
import { processFilesIntoText } from "@integrations/misc/extract-text"
import { showSystemNotification } from "@integrations/notifications"
import { ITerminalManager } from "@integrations/terminal/types"
import { startObservation } from "@langfuse/tracing"
import { BrowserSession } from "@services/browser/BrowserSession"
import { UrlContentFetcher } from "@services/browser/UrlContentFetcher"
import { featureFlagsService } from "@services/feature-flags"
@@ -135,6 +136,7 @@ export class Task {
private taskIsFavorited?: boolean
private cwd: string
private taskInitializationStartTime: number
private span
taskState: TaskState
@@ -332,6 +334,8 @@ export class Task {
updateTaskHistory: this.updateTaskHistory,
})
this.span = startObservation("task")
// Initialize context trackers
this.fileContextTracker = new FileContextTracker(controller, this.taskId)
this.modelContextTracker = new ModelContextTracker(this.taskId)
@@ -1324,6 +1328,7 @@ export class Task {
let includeFileDetails = true
while (!this.taskState.abort) {
const didEndLoop = await this.recursivelyMakeClineRequests(nextUserContent, includeFileDetails)
includeFileDetails = false // we only need file details the first time
// The way this agentic loop works is that cline will be given a task that he then calls tools to complete. unless there's an attempt_completion call, we keep responding back to him with his tool's responses until he either attempt_completion or does not use anymore tools. If he does not use anymore tools, we ask him to consider if he's completed the task and then call attempt_completion, otherwise proceed with completing the task.
@@ -1814,7 +1819,12 @@ export class Task {
}
// Response API requires native tool calls to be enabled
const stream = this.api.createMessage(systemPrompt, contextManagementMetadata.truncatedConversationHistory, tools)
const stream = this.api.createMessage(
systemPrompt,
contextManagementMetadata.truncatedConversationHistory,
tools,
this.span,
)
const iterator = stream[Symbol.asyncIterator]()
@@ -1967,6 +1977,7 @@ export class Task {
// Reset the automatic retry flag so the request can proceed
this.taskState.didAutomaticallyRetryFailedApiRequest = false
}
// delegate generator output from the recursive call
yield* this.attemptApiRequest(previousApiReqIndex)
return
@@ -2521,6 +2532,7 @@ export class Task {
taskMetrics.cacheWriteTokens += chunk.cacheWriteTokens ?? 0
taskMetrics.cacheReadTokens += chunk.cacheReadTokens ?? 0
taskMetrics.totalCost = chunk.totalCost ?? taskMetrics.totalCost
this.span.update({ output: assistantTextOnly })
break
case "reasoning": {
// Process the reasoning delta through the handler
@@ -2676,6 +2688,7 @@ export class Task {
await this.reinitExistingTaskFromId(this.taskId)
}
} finally {
this.span.end()
this.taskState.isStreaming = false
}
+3
View File
@@ -44,6 +44,7 @@ import { ExtensionRegistryInfo } from "./registry"
import { AuthService } from "./services/auth/AuthService"
import { LogoutReason } from "./services/auth/types"
import { telemetryService } from "./services/telemetry"
import { langfuse } from "./services/trace/langfuse"
import { SharedUriHandler } from "./services/uri/SharedUriHandler"
import { ShowMessageType } from "./shared/proto/host/window"
import { fileExistsAtPath } from "./utils/fs"
@@ -431,6 +432,8 @@ export async function activate(context: vscode.ExtensionContext) {
function setupHostProvider(context: ExtensionContext) {
console.log("Setting up vscode host providers...")
langfuse.start()
const createWebview = () => new VscodeWebviewProvider(context)
const createDiffView = () => new VscodeDiffViewProvider()
const createCommentReview = () => getVscodeCommentReviewController()
@@ -1,6 +1,6 @@
import { metrics } from "@opentelemetry/api"
import { logs } from "@opentelemetry/api-logs"
import { Resource } from "@opentelemetry/resources"
import { defaultResource, type Resource, resourceFromAttributes } from "@opentelemetry/resources"
import { BatchLogRecordProcessor, LoggerProvider } from "@opentelemetry/sdk-logs"
import { MeterProvider } from "@opentelemetry/sdk-metrics"
import { ATTR_SERVICE_NAME, ATTR_SERVICE_VERSION } from "@opentelemetry/semantic-conventions"
@@ -80,10 +80,12 @@ export class OpenTelemetryClientProvider {
}
// Create resource with service information
const resource = new Resource({
[ATTR_SERVICE_NAME]: "cline",
[ATTR_SERVICE_VERSION]: ExtensionRegistryInfo.version,
})
const resource = defaultResource().merge(
resourceFromAttributes({
[ATTR_SERVICE_NAME]: "cline",
[ATTR_SERVICE_VERSION]: ExtensionRegistryInfo.version,
}),
)
// Initialize metrics if configured
if (this.config.metricsExporter) {
+13
View File
@@ -0,0 +1,13 @@
import { LangfuseClient } from "@langfuse/client"
import { LangfuseSpanProcessor } from "@langfuse/otel"
import { NodeSDK } from "@opentelemetry/sdk-node"
// TODO: Convert this into a singleton pattern that can be configured with different providers apart from Langfuse
new LangfuseClient({
publicKey: process.env.LANGFUSE_PUBLIC_KEY,
secretKey: process.env.LANGFUSE_SECRET_KEY,
baseUrl: process.env.LANGFUSE_BASE_URL || "https://us.cloud.langfuse.com",
})
export const langfuse = new NodeSDK({ spanProcessors: [new LangfuseSpanProcessor()] })