feat: export AI SDK telemetry to local OTLP
This commit is contained in:
parent
4c4eef46f1
commit
c523f7ad2b
9 changed files with 331 additions and 121 deletions
|
|
@ -104,7 +104,14 @@
|
|||
"@ai-sdk/xai": "3.0.75",
|
||||
"@aws-sdk/credential-providers": "3.993.0",
|
||||
"@clack/prompts": "1.0.0-alpha.1",
|
||||
"@effect/opentelemetry": "catalog:",
|
||||
"@effect/platform-node": "catalog:",
|
||||
"@opentelemetry/api": "catalog:",
|
||||
"@opentelemetry/exporter-trace-otlp-http": "catalog:",
|
||||
"@opentelemetry/resources": "catalog:",
|
||||
"@opentelemetry/sdk-trace-base": "catalog:",
|
||||
"@opentelemetry/sdk-trace-node": "catalog:",
|
||||
"@opentelemetry/semantic-conventions": "catalog:",
|
||||
"@gitlab/opencode-gitlab-auth": "1.3.3",
|
||||
"@hono/node-server": "1.19.11",
|
||||
"@hono/node-ws": "1.3.0",
|
||||
|
|
|
|||
|
|
@ -21,7 +21,9 @@ import { Plugin } from "@/plugin"
|
|||
import { Skill } from "../skill"
|
||||
import { Effect, Context, Layer } from "effect"
|
||||
import { InstanceState } from "@/effect/instance-state"
|
||||
import { Observability } from "@/effect/oltp"
|
||||
import { makeRuntime } from "@/effect/run-service"
|
||||
import type { Tracer } from "@opentelemetry/api"
|
||||
|
||||
export namespace Agent {
|
||||
export const Info = z
|
||||
|
|
@ -345,38 +347,42 @@ export namespace Agent {
|
|||
const authInfo = yield* auth.get(model.providerID).pipe(Effect.orDie)
|
||||
const isOpenaiOauth = model.providerID === "openai" && authInfo?.type === "oauth"
|
||||
|
||||
const params = {
|
||||
experimental_telemetry: {
|
||||
isEnabled: cfg.experimental?.openTelemetry,
|
||||
metadata: {
|
||||
userId: cfg.username ?? "unknown",
|
||||
},
|
||||
},
|
||||
temperature: 0.3,
|
||||
messages: [
|
||||
...(isOpenaiOauth
|
||||
? []
|
||||
: system.map(
|
||||
(item): ModelMessage => ({
|
||||
role: "system",
|
||||
content: item,
|
||||
}),
|
||||
)),
|
||||
{
|
||||
role: "user",
|
||||
content: `Create an agent configuration based on this request: \"${input.description}\".\n\nIMPORTANT: The following identifiers already exist and must NOT be used: ${existing.map((i) => i.name).join(", ")}\n Return ONLY the JSON object, no other text, do not wrap in backticks`,
|
||||
},
|
||||
],
|
||||
model: language,
|
||||
schema: z.object({
|
||||
identifier: z.string(),
|
||||
whenToUse: z.string(),
|
||||
systemPrompt: z.string(),
|
||||
}),
|
||||
} satisfies Parameters<typeof generateObject>[0]
|
||||
const run = async (tracer: Tracer) => {
|
||||
const params = {
|
||||
experimental_telemetry: Observability.aiTelemetry({
|
||||
enabled: cfg.experimental?.openTelemetry,
|
||||
tracer,
|
||||
functionId: "Agent.generate",
|
||||
metadata: {
|
||||
userID: cfg.username ?? "unknown",
|
||||
providerID: resolved.providerID,
|
||||
modelID: resolved.id,
|
||||
},
|
||||
}),
|
||||
temperature: 0.3,
|
||||
messages: [
|
||||
...(isOpenaiOauth
|
||||
? []
|
||||
: system.map(
|
||||
(item): ModelMessage => ({
|
||||
role: "system",
|
||||
content: item,
|
||||
}),
|
||||
)),
|
||||
{
|
||||
role: "user",
|
||||
content: `Create an agent configuration based on this request: \"${input.description}\".\n\nIMPORTANT: The following identifiers already exist and must NOT be used: ${existing.map((i) => i.name).join(", ")}\n Return ONLY the JSON object, no other text, do not wrap in backticks`,
|
||||
},
|
||||
],
|
||||
model: language,
|
||||
schema: z.object({
|
||||
identifier: z.string(),
|
||||
whenToUse: z.string(),
|
||||
systemPrompt: z.string(),
|
||||
}),
|
||||
} satisfies Parameters<typeof generateObject>[0]
|
||||
|
||||
if (isOpenaiOauth) {
|
||||
return yield* Effect.promise(async () => {
|
||||
if (isOpenaiOauth) {
|
||||
const result = streamObject({
|
||||
...params,
|
||||
providerOptions: ProviderTransform.providerOptions(resolved, {
|
||||
|
|
@ -389,10 +395,12 @@ export namespace Agent {
|
|||
if (part.type === "error") throw part.error
|
||||
}
|
||||
return result.object
|
||||
})
|
||||
}
|
||||
|
||||
return generateObject(params).then((r) => r.object)
|
||||
}
|
||||
|
||||
return yield* Effect.promise(() => generateObject(params).then((r) => r.object))
|
||||
return yield* Observability.promise(run)
|
||||
}),
|
||||
})
|
||||
}),
|
||||
|
|
|
|||
|
|
@ -1,13 +1,45 @@
|
|||
import { Duration, Layer } from "effect"
|
||||
import * as NodeSdk from "@effect/opentelemetry/NodeSdk"
|
||||
import * as OtelResource from "@effect/opentelemetry/Resource"
|
||||
import * as OtelTracer from "@effect/opentelemetry/Tracer"
|
||||
import { context, trace, type AttributeValue, type Span, type Tracer } from "@opentelemetry/api"
|
||||
import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-http"
|
||||
import { BatchSpanProcessor } from "@opentelemetry/sdk-trace-base"
|
||||
import { Duration, Effect, Layer, ManagedRuntime, Option } from "effect"
|
||||
import * as Context from "effect/Context"
|
||||
import { FetchHttpClient } from "effect/unstable/http"
|
||||
import { Otlp } from "effect/unstable/observability"
|
||||
import { OtlpLogger, OtlpSerialization, OtlpTracer } from "effect/unstable/observability"
|
||||
import { normalizeServerUrl } from "@/account/url"
|
||||
import { EffectLogger } from "@/effect/logger"
|
||||
import { Flag } from "@/flag/flag"
|
||||
import { CHANNEL, VERSION } from "@/installation/meta"
|
||||
|
||||
export namespace Observability {
|
||||
export class AITracer extends Context.Service<AITracer, Tracer>()("@opencode/Observability/AITracer") {}
|
||||
|
||||
const clean = <T extends Record<string, unknown>>(value: T) =>
|
||||
Object.fromEntries(Object.entries(value).filter((entry) => entry[1] !== undefined)) as {
|
||||
[K in keyof T as undefined extends T[K] ? never : K]: Exclude<T[K], undefined>
|
||||
}
|
||||
|
||||
const parseHeaders = () =>
|
||||
Flag.OTEL_EXPORTER_OTLP_HEADERS
|
||||
? Flag.OTEL_EXPORTER_OTLP_HEADERS.split(",").reduce(
|
||||
(acc, item) => {
|
||||
const at = item.indexOf("=")
|
||||
if (at < 1 || at === item.length - 1) return acc
|
||||
acc[item.slice(0, at)] = item.slice(at + 1)
|
||||
return acc
|
||||
},
|
||||
{} as Record<string, string>,
|
||||
)
|
||||
: undefined
|
||||
|
||||
const base = Flag.OTEL_EXPORTER_OTLP_ENDPOINT
|
||||
export const enabled = !!base
|
||||
const root = base ? normalizeServerUrl(base) : undefined
|
||||
const traces = Flag.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT ?? (root ? `${root}/v1/traces` : undefined)
|
||||
const logs = Flag.OTEL_EXPORTER_OTLP_LOGS_ENDPOINT ?? (root ? `${root}/v1/logs` : undefined)
|
||||
|
||||
export const enabled = !!traces || !!logs
|
||||
|
||||
const resource = {
|
||||
serviceName: "opencode",
|
||||
|
|
@ -18,24 +50,86 @@ export namespace Observability {
|
|||
},
|
||||
}
|
||||
|
||||
const headers = Flag.OTEL_EXPORTER_OTLP_HEADERS
|
||||
? Flag.OTEL_EXPORTER_OTLP_HEADERS.split(",").reduce(
|
||||
(acc, x) => {
|
||||
const [key, value] = x.split("=")
|
||||
acc[key] = value
|
||||
return acc
|
||||
},
|
||||
{} as Record<string, string>,
|
||||
)
|
||||
: undefined
|
||||
const headers = parseHeaders()
|
||||
|
||||
export const layer = !base
|
||||
? EffectLogger.layer
|
||||
: Otlp.layerJson({
|
||||
baseUrl: base,
|
||||
loggerExportInterval: Duration.seconds(1),
|
||||
loggerMergeWithExisting: true,
|
||||
const tracer = traces
|
||||
? OtlpTracer.layer({
|
||||
url: traces,
|
||||
resource,
|
||||
headers,
|
||||
}).pipe(Layer.provide(EffectLogger.layer), Layer.provide(FetchHttpClient.layer))
|
||||
})
|
||||
: Layer.empty
|
||||
|
||||
const logger = logs
|
||||
? OtlpLogger.layer({
|
||||
url: logs,
|
||||
resource,
|
||||
headers,
|
||||
exportInterval: Duration.seconds(1),
|
||||
mergeWithExisting: true,
|
||||
})
|
||||
: Layer.empty
|
||||
|
||||
const ai = traces
|
||||
? Layer.effect(AITracer, Effect.service(OtelTracer.OtelTracer)).pipe(
|
||||
Layer.provide(
|
||||
OtelTracer.layerTracer.pipe(
|
||||
Layer.provide(
|
||||
NodeSdk.layerTracerProvider(new BatchSpanProcessor(new OTLPTraceExporter({ url: traces, headers }))),
|
||||
),
|
||||
Layer.provide(OtelResource.layer(resource)),
|
||||
),
|
||||
),
|
||||
)
|
||||
: Layer.succeed(AITracer, trace.getTracer(resource.serviceName, resource.serviceVersion))
|
||||
|
||||
export const layer =
|
||||
!traces && !logs
|
||||
? Layer.mergeAll(EffectLogger.layer, ai)
|
||||
: Layer.mergeAll(tracer, logger, ai).pipe(
|
||||
Layer.provide(EffectLogger.layer),
|
||||
Layer.provide(OtlpSerialization.layerJson),
|
||||
Layer.provide(FetchHttpClient.layer),
|
||||
)
|
||||
|
||||
const runtime = ManagedRuntime.make(layer)
|
||||
const aiRuntime = ManagedRuntime.make(ai)
|
||||
|
||||
const withSpan = <A>(span: Option.Option<Span>, fn: () => A): A =>
|
||||
Option.match(span, {
|
||||
onNone: fn,
|
||||
onSome: (span) => context.with(trace.setSpan(context.active(), span), fn),
|
||||
})
|
||||
|
||||
const withActiveParent = <A, E, R>(effect: Effect.Effect<A, E, R>) => {
|
||||
const active = trace.getActiveSpan()
|
||||
if (!active) return effect
|
||||
return effect.pipe(OtelTracer.withSpanContext(active.spanContext()))
|
||||
}
|
||||
|
||||
export const runPromise = <A, E>(effect: Effect.Effect<A, E>) => runtime.runPromise(withActiveParent(effect))
|
||||
|
||||
export const runFork = <A, E>(effect: Effect.Effect<A, E>) => runtime.runFork(withActiveParent(effect))
|
||||
|
||||
export const promise = <A>(fn: (tracer: Tracer) => Promise<A> | A) =>
|
||||
Effect.gen(function* () {
|
||||
const span = yield* Effect.option(OtelTracer.currentOtelSpan)
|
||||
const tracer = yield* Effect.promise(() => aiRuntime.runPromise(Effect.service(AITracer)))
|
||||
return yield* Effect.promise(() => Promise.resolve(withSpan(span, () => fn(tracer))))
|
||||
})
|
||||
|
||||
export const aiTelemetry = (input: {
|
||||
enabled: boolean | undefined
|
||||
tracer: Tracer
|
||||
functionId: string
|
||||
metadata?: Record<string, AttributeValue | undefined>
|
||||
}) => {
|
||||
if (!input.enabled || !traces) return { isEnabled: false as const }
|
||||
return {
|
||||
isEnabled: true as const,
|
||||
functionId: input.functionId,
|
||||
tracer: input.tracer,
|
||||
metadata: input.metadata ? clean(input.metadata) : undefined,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -12,6 +12,8 @@ function falsy(key: string) {
|
|||
|
||||
export namespace Flag {
|
||||
export const OTEL_EXPORTER_OTLP_ENDPOINT = process.env["OTEL_EXPORTER_OTLP_ENDPOINT"]
|
||||
export const OTEL_EXPORTER_OTLP_TRACES_ENDPOINT = process.env["OTEL_EXPORTER_OTLP_TRACES_ENDPOINT"]
|
||||
export const OTEL_EXPORTER_OTLP_LOGS_ENDPOINT = process.env["OTEL_EXPORTER_OTLP_LOGS_ENDPOINT"]
|
||||
export const OTEL_EXPORTER_OTLP_HEADERS = process.env["OTEL_EXPORTER_OTLP_HEADERS"]
|
||||
|
||||
export const OPENCODE_AUTO_SHARE = truthy("OPENCODE_AUTO_SHARE")
|
||||
|
|
|
|||
|
|
@ -21,67 +21,14 @@ import { Wildcard } from "@/util/wildcard"
|
|||
import { SessionID } from "@/session/schema"
|
||||
import { Auth } from "@/auth"
|
||||
import { Installation } from "@/installation"
|
||||
import { Observability } from "@/effect/oltp"
|
||||
import type { Tracer } from "@opentelemetry/api"
|
||||
|
||||
export namespace LLM {
|
||||
const log = Log.create({ service: "llm" })
|
||||
export const OUTPUT_TOKEN_MAX = ProviderTransform.OUTPUT_TOKEN_MAX
|
||||
|
||||
export type StreamInput = {
|
||||
user: MessageV2.User
|
||||
sessionID: string
|
||||
parentSessionID?: string
|
||||
model: Provider.Model
|
||||
agent: Agent.Info
|
||||
permission?: Permission.Ruleset
|
||||
system: string[]
|
||||
messages: ModelMessage[]
|
||||
small?: boolean
|
||||
tools: Record<string, Tool>
|
||||
retries?: number
|
||||
toolChoice?: "auto" | "required" | "none"
|
||||
}
|
||||
|
||||
export type StreamRequest = StreamInput & {
|
||||
abort: AbortSignal
|
||||
}
|
||||
|
||||
export type Event = Awaited<ReturnType<typeof stream>>["fullStream"] extends AsyncIterable<infer T> ? T : never
|
||||
|
||||
export interface Interface {
|
||||
readonly stream: (input: StreamInput) => Stream.Stream<Event, unknown>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/LLM") {}
|
||||
|
||||
export const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
return Service.of({
|
||||
stream(input) {
|
||||
return Stream.scoped(
|
||||
Stream.unwrap(
|
||||
Effect.gen(function* () {
|
||||
const ctrl = yield* Effect.acquireRelease(
|
||||
Effect.sync(() => new AbortController()),
|
||||
(ctrl) => Effect.sync(() => ctrl.abort()),
|
||||
)
|
||||
|
||||
const result = yield* Effect.promise(() => LLM.stream({ ...input, abort: ctrl.signal }))
|
||||
|
||||
return Stream.fromAsyncIterable(result.fullStream, (e) =>
|
||||
e instanceof Error ? e : new Error(String(e)),
|
||||
)
|
||||
}),
|
||||
),
|
||||
)
|
||||
},
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
export const defaultLayer = layer
|
||||
|
||||
export async function stream(input: StreamRequest) {
|
||||
const request = async (input: StreamRequest, tracer: Tracer) => {
|
||||
const l = log
|
||||
.clone()
|
||||
.tag("providerID", input.model.providerID)
|
||||
|
|
@ -384,16 +331,82 @@ export namespace LLM {
|
|||
},
|
||||
],
|
||||
}),
|
||||
experimental_telemetry: {
|
||||
isEnabled: cfg.experimental?.openTelemetry,
|
||||
experimental_telemetry: Observability.aiTelemetry({
|
||||
enabled: cfg.experimental?.openTelemetry,
|
||||
tracer,
|
||||
functionId: "LLM.stream",
|
||||
metadata: {
|
||||
userId: cfg.username ?? "unknown",
|
||||
sessionId: input.sessionID,
|
||||
userID: cfg.username ?? "unknown",
|
||||
sessionID: input.sessionID,
|
||||
providerID: input.model.providerID,
|
||||
modelID: input.model.id,
|
||||
agent: input.agent.name,
|
||||
},
|
||||
},
|
||||
}),
|
||||
})
|
||||
}
|
||||
|
||||
export type StreamInput = {
|
||||
user: MessageV2.User
|
||||
sessionID: string
|
||||
parentSessionID?: string
|
||||
model: Provider.Model
|
||||
agent: Agent.Info
|
||||
permission?: Permission.Ruleset
|
||||
system: string[]
|
||||
messages: ModelMessage[]
|
||||
small?: boolean
|
||||
tools: Record<string, Tool>
|
||||
retries?: number
|
||||
toolChoice?: "auto" | "required" | "none"
|
||||
}
|
||||
|
||||
export type StreamRequest = StreamInput & {
|
||||
abort: AbortSignal
|
||||
}
|
||||
|
||||
export type Event = Awaited<ReturnType<typeof stream>>["fullStream"] extends AsyncIterable<infer T> ? T : never
|
||||
|
||||
export interface Interface {
|
||||
readonly stream: (input: StreamInput) => Stream.Stream<Event, unknown>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/LLM") {}
|
||||
|
||||
export const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
return Service.of({
|
||||
stream(input) {
|
||||
return Stream.scoped(
|
||||
Stream.unwrap(
|
||||
Effect.gen(function* () {
|
||||
const ctrl = yield* Effect.acquireRelease(
|
||||
Effect.sync(() => new AbortController()),
|
||||
(ctrl) => Effect.sync(() => ctrl.abort()),
|
||||
)
|
||||
|
||||
const result = yield* Observability.promise((tracer) =>
|
||||
request({ ...input, abort: ctrl.signal }, tracer),
|
||||
)
|
||||
|
||||
return Stream.fromAsyncIterable(result.fullStream, (e) =>
|
||||
e instanceof Error ? e : new Error(String(e)),
|
||||
)
|
||||
}),
|
||||
),
|
||||
)
|
||||
},
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
export const defaultLayer = layer
|
||||
|
||||
export async function stream(input: StreamRequest) {
|
||||
return Observability.runPromise(Observability.promise((tracer) => request(input, tracer)))
|
||||
}
|
||||
|
||||
function resolveTools(input: Pick<StreamInput, "tools" | "agent" | "permission" | "user">) {
|
||||
const disabled = Permission.disabled(
|
||||
Object.keys(input.tools),
|
||||
|
|
|
|||
|
|
@ -45,6 +45,7 @@ import { decodeDataUrl } from "@/util/data-url"
|
|||
import { Process } from "@/util/process"
|
||||
import { Cause, Effect, Exit, Layer, Option, Scope, Context } from "effect"
|
||||
import { EffectLogger } from "@/effect/logger"
|
||||
import { Observability } from "@/effect/oltp"
|
||||
import { InstanceState } from "@/effect/instance-state"
|
||||
import { makeRuntime } from "@/effect/run-service"
|
||||
import { TaskTool, type TaskPromptOps } from "@/tool/task"
|
||||
|
|
@ -106,9 +107,8 @@ export namespace SessionPrompt {
|
|||
const llm = yield* LLM.Service
|
||||
|
||||
const run = {
|
||||
promise: <A, E>(effect: Effect.Effect<A, E>) =>
|
||||
Effect.runPromise(effect.pipe(Effect.provide(EffectLogger.layer))),
|
||||
fork: <A, E>(effect: Effect.Effect<A, E>) => Effect.runFork(effect.pipe(Effect.provide(EffectLogger.layer))),
|
||||
promise: <A, E>(effect: Effect.Effect<A, E>) => Observability.runPromise(effect),
|
||||
fork: <A, E>(effect: Effect.Effect<A, E>) => Observability.runFork(effect),
|
||||
}
|
||||
|
||||
const cancel = Effect.fn("SessionPrompt.cancel")(function* (sessionID: SessionID) {
|
||||
|
|
|
|||
|
|
@ -75,8 +75,17 @@ export namespace Tool {
|
|||
Effect.gen(function* () {
|
||||
const toolInfo = init instanceof Function ? { ...(yield* init()) } : { ...init }
|
||||
const execute = toolInfo.execute
|
||||
toolInfo.execute = (args, ctx) =>
|
||||
Effect.gen(function* () {
|
||||
toolInfo.execute = (args, ctx) => {
|
||||
const ann = Object.fromEntries(
|
||||
Object.entries({
|
||||
tool: id,
|
||||
agent: ctx.agent,
|
||||
sessionID: ctx.sessionID,
|
||||
messageID: ctx.messageID,
|
||||
callID: ctx.callID,
|
||||
}).filter((entry) => entry[1] !== undefined),
|
||||
)
|
||||
return Effect.gen(function* () {
|
||||
yield* Effect.try({
|
||||
try: () => toolInfo.parameters.parse(args),
|
||||
catch: (error) => {
|
||||
|
|
@ -104,7 +113,14 @@ export namespace Tool {
|
|||
...(truncated.truncated && { outputPath: truncated.outputPath }),
|
||||
},
|
||||
}
|
||||
}).pipe(Effect.orDie)
|
||||
}).pipe(
|
||||
Effect.annotateLogs(ann),
|
||||
Effect.annotateSpans(ann),
|
||||
Effect.withLogSpan(`Tool.${id}`),
|
||||
Effect.withSpan(`Tool.${id}`),
|
||||
Effect.orDie,
|
||||
)
|
||||
}
|
||||
return toolInfo
|
||||
})
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue