Documentation
¶
Overview ¶
Package otel wires weft's observability to one or more OpenTelemetry destinations at once — the local sink, Weft Studio, Datadog, Langfuse, any OTLP endpoint or your own exporters — each with its own signals and content policy (ADR 0024, S2).
One Install, one set of spans and records, fanned out:
defer otel.Install(
otel.Local("weft.db"), // local sink: offline, replay-grade, content on
otel.Studio("https://studio.example", token), // content on, traces + logs
otel.Datadog(), // the local Datadog Agent's OTLP intake, content off
otel.Langfuse("https://cloud.langfuse.com", pk, sk), // traces only
)()
The rules that shape it:
- Content is captured once and shaped per destination [D2]. Every processor registered here answers the Logs API's Enabled by event name: the core asks one question ("weft.messages") and emits content when any destination wants it. Content-off chains clone the record first (the SDK hands every processor the same pointer), strip the body with weft.StripContent and drop messages records; content-on chains redact and cap event and delta bodies and redact (never cap) the transcript records part by part. The core reads no environment variable — this module does, here only.
- The run tracker keeps the set of open runs and emits one weft.heartbeat record per open run per interval, so a 45 s tool call reads running and a crashed process stops heartbeating (the "interrupted" rule stays honest). Sinks never store heartbeats.
- The Local destination writes obsdb/sqlite synchronously: nothing emitted is lost, and the DB's hub is setup A's live lane.
- Install never panics and never fails the program; Start returns errors and leaves the globals alone under NoGlobal (tests).
This module depends on the OTel SDK and the OTLP/HTTP exporters; the core never does. It never imports weft/thread (the runtime link is its own module, weft/runtime, D6).
Index ¶
- func FromSDKRecords(recs []sdklog.Record) []obsdb.Record
- func FromSDKSpans(spans []sdktrace.ReadOnlySpan) []obsdb.Span
- func Install(opts ...Option) func()
- func LocalDB() obsdb.DB
- func StudioEndpoint() (string, string)
- type ContentConfig
- type DestOption
- func BatchDelay(d time.Duration) DestOption
- func DatadogEndpoint(url string) DestOption
- func Headers(h map[string]string) DestOption
- func Insecure() DestOption
- func NoContent() DestOption
- func NoDeltas() DestOption
- func Signals(traces, logs bool) DestOption
- func Timeout(d time.Duration) DestOption
- func WithContent(cfg ...ContentConfig) DestOption
- type Option
- func Content(cfg ContentConfig) Option
- func Datadog(opts ...DestOption) Option
- func Exporters(spans sdktrace.SpanExporter, logs sdklog.Exporter, opts ...DestOption) Option
- func Heartbeat(every time.Duration) Option
- func Langfuse(host, publicKey, secretKey string, opts ...DestOption) Option
- func Local(path string, opts ...DestOption) Option
- func NoEnv() Option
- func NoGlobal() Option
- func OTLP(endpoint string, opts ...DestOption) Option
- func Resource(r *sdkresource.Resource) Option
- func Sampler(s sdktrace.Sampler) Option
- func Service(name string) Option
- func Studio(url, token string, opts ...DestOption) Option
- type Pipeline
- func (p *Pipeline) ForceFlush(ctx context.Context) error
- func (p *Pipeline) LocalDB() obsdb.DB
- func (p *Pipeline) LoggerProvider() *sdklog.LoggerProvider
- func (p *Pipeline) Shutdown(ctx context.Context) error
- func (p *Pipeline) StudioEndpoint() (string, string)
- func (p *Pipeline) TracerProvider() *sdktrace.TracerProvider
Examples ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func FromSDKRecords ¶
FromSDKRecords converts SDK log records (what a log Exporter receives) into the obsdb model, matching obsdb.FromOTLPLogs value for value on the same data.
func FromSDKSpans ¶
func FromSDKSpans(spans []sdktrace.ReadOnlySpan) []obsdb.Span
FromSDKSpans converts SDK spans (what a SpanExporter receives) into the obsdb model. It lives in this module because it needs the SDK; a package under obsdb would put the SDK in obsdb's go.mod (S3.3). For the same data it produces exactly the values obsdb.FromOTLPTraces does — the golden test pins that against obsdb's OTLP fixtures.
func Install ¶
func Install(opts ...Option) func()
Install starts the pipeline and registers it globally. It never panics and never fails the program: with no options at all (and no environment destinations) it writes the local sink only (the zero-config rule), and a destination that cannot be built is logged (slog, WARN) and skipped. When no destination option was passed and every environment destination failed to build, it falls back to the local sink with one WARN rather than recording nothing (Start returns the error instead); a program that passed its own destinations gets no destination it did not name. The returned function flushes and shuts everything down; call it on exit.
Example ¶
Install is the whole setup: it registers the OTel globals the core reads, never fails the program, and returns the shutdown to call on exit. NoEnv keeps the example hermetic — no WEFT_* / OTEL_* variable can add a destination behind it.
package main
import (
"context"
"fmt"
"log"
"os"
"path/filepath"
"github.com/weftgo/weft"
"github.com/weftgo/weft/obsdb"
"github.com/weftgo/weft/obsdb/sqlite"
"github.com/weftgo/weft/otel"
"github.com/weftgo/weft/wefttest"
)
func main() {
dir, err := os.MkdirTemp("", "weft-otel-example-")
if err != nil {
log.Fatal(err)
}
defer func() { _ = os.RemoveAll(dir) }()
shutdown := otel.Install(
otel.NoEnv(),
otel.Local(filepath.Join(dir, "weft.db")),
)
// A scripted run through the installed globals — a real program
// wires nothing: the core reads the providers Install registered.
agt := weft.New(
wefttest.Script(wefttest.Say("all done")),
weft.Name("example"),
)
if _, err := agt.Generate(context.Background(), weft.Prompt("go"),
weft.Metadata(map[string]string{"weft.session.id": "example-session"})); err != nil {
log.Fatal(err)
}
shutdown() // flush every destination, close the local DB last
// The sink reopens as a plain obsdb.DB: the run row with its
// derived status, and the transcript the messages records hold.
db, err := sqlite.Open(filepath.Join(dir, "weft.db"))
if err != nil {
log.Fatal(err)
}
defer func() { _ = db.Close() }()
page, err := db.Runs(context.Background(), obsdb.RunQuery{})
if err != nil {
log.Fatal(err)
}
fmt.Println(page.Runs[0].Status, page.Runs[0].MessageCount)
}
Output: succeeded 2
func LocalDB ¶
LocalDB returns the installed pipeline's local obsdb.DB — for studio.DB(otel.LocalDB()) (setup A, [D4]). nil before Install, after its shutdown, or without a Local destination.
S2.1 names this otel.Local(), which cannot coexist with the Local(path) destination option in Go; LocalDB is the recorded deviation (notes-lane-a2.md).
func StudioEndpoint ¶
StudioEndpoint returns the installed pipeline's Studio destination, which weft/runtime dials by default. "" without one (or after the pipeline's shutdown).
Types ¶
type ContentConfig ¶
type ContentConfig struct {
// MaxBytes caps one event field or delta body; 0 = 32 KiB; negative
// = unlimited. Never applied to messages records — a capped
// transcript is not replay-grade.
MaxBytes int
// Redact is applied to a content field before the cap; nil = identity.
// It sees event and delta bodies and weft.messages transcript
// records, the latter part by part with the kind the same content
// has on the event path: text parts (every role) weft.ContentText,
// reasoning weft.ContentReasoning, tool-call args weft.ContentArgs (an
// output that is not JSON is carried as a JSON string), tool results
// weft.ContentResult. Ids, names, roles, signatures and file parts are
// not passed to it; a transcript record is redacted, never capped. It
// runs on the run's goroutine: a panic in it is contained — the event
// goes out stripped, the transcript record not at all — and counted
// as the destination's drop; so is a transcript record that cannot be
// decoded for redaction. Nothing it was given is ever sent unredacted.
Redact func(kind weft.ContentKind, s string) string
}
ContentConfig carries the caps and redaction that live per destination, not in the core ([D2]: the core cannot discover them — the OTel global never hands out the SDK provider).
type DestOption ¶
type DestOption interface {
// contains filtered or unexported methods
}
DestOption configures one destination.
func BatchDelay ¶
func BatchDelay(d time.Duration) DestOption
BatchDelay sets the destination's batch flush interval (Studio defaults logs 200 ms / spans 1 s; the others the SDK's defaults).
func DatadogEndpoint ¶
func DatadogEndpoint(url string) DestOption
DatadogEndpoint overrides the Datadog destination's OTLP intake URL (default http://localhost:4318).
func Headers ¶
func Headers(h map[string]string) DestOption
Headers adds HTTP headers to the destination's requests (OTLP-family destinations).
func Insecure ¶
func Insecure() DestOption
Insecure allows http:// to a non-loopback host (OTLP-family destinations). Loopback http is always allowed, and so is the environment's OTEL_EXPORTER_OTLP_ENDPOINT written with http:// (the operator's opt-in; see envDestinations) — never WEFT_STUDIO_URL.
func NoContent ¶
func NoContent() DestOption
NoContent turns content off: the chain strips event and delta bodies (weft.StripContent) and drops messages records before its exporter.
func NoDeltas ¶
func NoDeltas() DestOption
NoDeltas drops text/reasoning/args deltas before export, to save bandwidth when live streaming is not wanted. Deltas are never stored by the sinks either way (Q4); the core's durable sequence is untouched (they number on their own counter, D3).
func Signals ¶
func Signals(traces, logs bool) DestOption
Signals selects which signals this destination takes. The defaults are per destination (Langfuse is traces only; the rest both).
func Timeout ¶
func Timeout(d time.Duration) DestOption
Timeout sets the destination's export request timeout.
func WithContent ¶
func WithContent(cfg ...ContentConfig) DestOption
WithContent turns content on for this destination, optionally overriding the pipeline's ContentConfig.
type Option ¶
type Option interface {
// contains filtered or unexported methods
}
Option configures the pipeline (Install/Start).
func Content ¶
func Content(cfg ContentConfig) Option
Content sets the pipeline-wide default for destinations with content on; a destination's WithContent overrides it.
func Datadog ¶
func Datadog(opts ...DestOption) Option
Datadog exports to the local Datadog Agent's OTLP intake (http://localhost:4318; the Agent binds the standard OTLP ports when its OTLP receiver is enabled). Content off, traces + logs, SDK default batching. DatadogEndpoint overrides the URL.
func Exporters ¶
func Exporters(spans sdktrace.SpanExporter, logs sdklog.Exporter, opts ...DestOption) Option
Exporters wires your own exporters (vendor-specific, stdout, tests): whichever is non-nil gets a chain. Content off, SDK default batching.
func Heartbeat ¶
Heartbeat sets how often the run tracker emits one weft.heartbeat record per open run (the default 10s; 0 or less disables). Three missed intervals are what obsdb's InterruptedAfter reads as interrupted.
func Langfuse ¶
func Langfuse(host, publicKey, secretKey string, opts ...DestOption) Option
Langfuse exports traces to Langfuse's OTLP endpoint (<host>/api/public/otel/v1/traces, Basic auth base64(pk:sk); OTLP/HTTP only — gRPC is not accepted). Traces only: weft's spans carry no content (ADR 0016 O7) and Langfuse ingests traces, so it gets timing, usage, tool names and the identity chain; span-level content is the post-v1 follow-up F1.
func Local ¶
func Local(path string, opts ...DestOption) Option
Local is the local sink: obsdb/sqlite at path ("" = $WEFT_DB or ./.weft/weft.db), written synchronously — nothing emitted is lost, and the DB's in-process hub is the live lane setup A shares with Studio. Content on, traces + logs.
func NoEnv ¶
func NoEnv() Option
NoEnv ignores the WEFT_* / OTEL_* destination variables. Explicit options are unaffected.
func NoGlobal ¶
func NoGlobal() Option
NoGlobal leaves the OTel globals alone (Start only; Install always registers). For tests and dependency-injected programs.
func OTLP ¶
func OTLP(endpoint string, opts ...DestOption) Option
OTLP exports to any OTLP/HTTP endpoint (http(s)://host:port). Content off, traces + logs, SDK default batching; Headers adds request headers.
func Resource ¶
func Resource(r *sdkresource.Resource) Option
Resource merges r over the detected resource.
func Sampler ¶
Sampler sets the trace sampler (default ParentBased(AlwaysOn)). Records are log records; sampling never applies to them.
func Studio ¶
func Studio(url, token string, opts ...DestOption) Option
Studio exports OTLP/HTTP protobuf to a Weft Studio with a bearer token. Content on, traces + logs, short batches (logs 200 ms, spans 1 s).
type Pipeline ¶
type Pipeline struct {
// contains filtered or unexported fields
}
Pipeline is the running fan-out: one TracerProvider and one LoggerProvider, one processor chain per destination, one run tracker. Create it with Start (or Install, which never fails).
func Start ¶
Start builds and starts the pipeline: one provider pair, the run tracker, one processor chain per destination (explicit options and environment combined, de-duplicated with the explicit destination winning). It fails if any destination fails to build — what was already built is shut down again, exporters handed to Exporters included; NoGlobal leaves the OTel globals alone (tests).
func (*Pipeline) ForceFlush ¶
ForceFlush flushes every destination, all at once: each destination's span and log chains flush on their own goroutines under ctx, so a slow or hung destination costs only its own data, never the budget of the one behind it (flushed through the providers they go one after another).
func (*Pipeline) LocalDB ¶
LocalDB returns the pipeline's local obsdb.DB — the handle setup A shares with Studio (studio.DB(otel.LocalDB())). nil without a Local destination.
The spec's S2.1 names this accessor otel.Local(), which cannot coexist with the Local(path) destination option in Go (one package, one name); LocalDB is the recorded deviation — see weft-otel-build/notes-lane-a2.md.
func (*Pipeline) LoggerProvider ¶
func (p *Pipeline) LoggerProvider() *sdklog.LoggerProvider
LoggerProvider returns the pipeline's logger provider.
func (*Pipeline) Shutdown ¶
Shutdown runs S2.4's order: stop heartbeats → flush every destination in parallel (5 s budget) → shut down the providers → close the local DB last.
func (*Pipeline) StudioEndpoint ¶
StudioEndpoint returns the pipeline's Studio destination (url and token), which weft/runtime dials by default. "" without one.
func (*Pipeline) TracerProvider ¶
func (p *Pipeline) TracerProvider() *sdktrace.TracerProvider
TracerProvider returns the pipeline's tracer provider.