otel

package module
v0.2.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 7, 2026 License: MIT Imports: 28 Imported by: 2

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

Examples

Constants

This section is empty.

Variables

This section is empty.

Functions

func FromSDKRecords

func FromSDKRecords(recs []sdklog.Record) []obsdb.Record

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

func LocalDB() obsdb.DB

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

func StudioEndpoint() (string, string)

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

func Heartbeat(every time.Duration) Option

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

func Sampler(s sdktrace.Sampler) Option

Sampler sets the trace sampler (default ParentBased(AlwaysOn)). Records are log records; sampling never applies to them.

func Service

func Service(name string) Option

Service sets resource service.name. Else OTEL_SERVICE_NAME, else the binary's name.

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

func Start(ctx context.Context, opts ...Option) (*Pipeline, error)

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

func (p *Pipeline) ForceFlush(ctx context.Context) error

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

func (p *Pipeline) LocalDB() obsdb.DB

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

func (p *Pipeline) Shutdown(ctx context.Context) error

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

func (p *Pipeline) StudioEndpoint() (string, string)

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.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL