Documentation
¶
Overview ¶
Package mw holds the reference middleware for weft's two seams: model middleware (weft.WrapModel) and tool middleware (weft.WrapTools). Each is a small, dependency-free value that shows the shape of its seam; copy it, compose it, or write your own.
Model seam — around every model call the loop makes:
weft.New(model, weft.WrapModel(
mw.Log(logger), // outermost: sees retries and fallbacks
mw.Fallback(backupModel), // fail over once Retry has given up
mw.Retry(mw.MaxRetries(3)), // transient failures, backoff, retry-after
mw.RepairJSON(), // innermost: fixes truncated tool-call args
))
Order matters: Fallback outside Retry retries the primary to exhaustion before switching models (retry, then fail over). The reverse — Retry outside Fallback — retries the fallback chain as a whole, so the primary gets one attempt per retry cycle (up to MaxRetries+1 in total, each followed by the backup) and Retry never sees its transient failures.
Tool seam — around every tool call the loop dispatches:
weft.New(model, weft.WrapTools(
mw.Audit(logger), // observation: every call, with its cause
mw.Allow(policy.Permits), // decision: DENIED results the model sees
mw.MapErrors(nil), // shaping: plain errors → INTERNAL codes
), tools...)
Middleware that verifies something (a user, a tenant, a quota) adds it to ctx before calling next, and tools read it back through a typed accessor — the "typed request context" convention documented in docs/life-of-a-call.md and weft's ExampleWrapTools_context.
Model middleware should forward the inner model's identity; wefttest.ConformInfoT (in the root module's wefttest, not here — the reference middleware and its checker stay mutually discoverable) turns that convention into a checked fact for your own middleware.
Example (PiiScrubMiddleware) ¶
A PII scrubber sits on the tool seam and masks what the model — and so every transcript downstream — is allowed to see. Both channels are scrubbed: the result text, and the error's text too — the loop renders err.Error() into the transcript verbatim, so scrubbing only the success path would leak through every failure. Rewriting the message means dropping the cause chain on purpose; a scrubber that keeps the original reachable is no scrubber.
package main
import (
"context"
"errors"
"fmt"
"regexp"
"github.com/weftgo/weft"
"github.com/weftgo/weft/wefttest"
)
func main() {
email := regexp.MustCompile(`[a-z0-9._%+-]+@[a-z0-9.-]+\.[a-z]{2,}`)
scrub := func(next weft.ToolCaller) weft.ToolCaller {
return func(ctx context.Context, call weft.ToolCallPart) (string, error) {
out, err := next(ctx, call)
if err != nil {
return "", errors.New(email.ReplaceAllString(err.Error(), "[redacted]"))
}
return email.ReplaceAllString(out, "[redacted]"), nil
}
}
lookup := weft.Tool("lookup", "", func(_ context.Context, _ struct{}) (string, error) {
return `{"email":"wajih@example.com","status":"delivered"}`, nil
})
agt := weft.New(wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "lookup"}),
wefttest.Say("done"),
), lookup, weft.WrapTools(scrub))
res, err := agt.Generate(context.Background(), weft.Prompt("look up my order"))
if err != nil {
fmt.Println(err)
return
}
for _, m := range res.Messages {
if m.Role != weft.RoleTool {
continue
}
for _, p := range m.Content {
if r, ok := p.(weft.ToolResultPart); ok {
fmt.Println(r.Content)
}
}
}
}
Output: {"email":"[redacted]","status":"delivered"}
Example (RateLimitMiddleware) ¶
Rate limiting is middleware, not core — this one spaces calls out with nothing but the standard library (the reason mw ships no RateLimit: golang.org/x/time would be the module's first dependency).
package main
import (
"context"
"fmt"
"log"
"sync"
"time"
"github.com/weftgo/weft"
"github.com/weftgo/weft/wefttest"
)
func main() {
const interval = 20 * time.Millisecond
var mu sync.Mutex
nextAt := time.Now()
reserve := func() time.Duration {
mu.Lock()
defer mu.Unlock()
wait := time.Until(nextAt)
if wait < 0 {
wait = 0
}
nextAt = nextAt.Add(interval)
return wait
}
limiter := func(next weft.ToolCaller) weft.ToolCaller {
return func(ctx context.Context, call weft.ToolCallPart) (string, error) {
if wait := reserve(); wait > 0 {
t := time.NewTimer(wait)
defer t.Stop()
select {
case <-t.C:
case <-ctx.Done():
return "", ctx.Err()
}
}
return next(ctx, call)
}
}
ping := weft.Tool("ping", "", func(_ context.Context, _ struct{}) (string, error) { return "pong", nil })
agt := weft.New(wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "ping"}, wefttest.Call{Name: "ping"}, wefttest.Call{Name: "ping"}),
wefttest.Say("done"),
), ping, weft.WrapTools(limiter))
start := time.Now()
res, err := agt.Generate(context.Background(), weft.Prompt("x"))
if err != nil {
log.Fatal(err)
}
fmt.Println(res.Text(), "spread over", time.Since(start) >= 2*interval)
}
Output: done spread over true
Example (ResponseCacheMiddleware) ¶
A response cache is model middleware keyed on the ModelRequest: the same prompt and catalogue replay the recorded answer without a provider call. Cache invalidation is the caller's policy — here the whole request is the key, so any change misses. The map carries a mutex because a Model must survive the Agent's concurrent reuse (the rate-limit example's rule; an unsynchronized map races the moment two runs share the cache).
package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"iter"
"sync"
"github.com/weftgo/weft"
"github.com/weftgo/weft/wefttest"
)
func main() {
key := func(req weft.ModelRequest) string {
b, _ := json.Marshal(struct {
System string
Messages []weft.Message
Tools []string
}{req.System, req.Messages, toolNamesOf(req.Tools)})
sum := sha256.Sum256(b)
return hex.EncodeToString(sum[:8])
}
shared := &responseCache{cache: map[string][]weft.ModelEvent{}}
model := wefttest.Script(wefttest.Say("fresh"), wefttest.Say("fresh"))
agt := weft.New(model, weft.WrapModel(func(next weft.Model) weft.Model {
return cacheModel{next: next, c: shared, key: key}
}))
for range 3 {
if _, err := agt.Generate(context.Background(), weft.Prompt("same question")); err != nil {
fmt.Println(err)
return
}
}
fmt.Printf("provider calls made: %d of 3\n", 3-shared.Hits())
}
// responseCache is the shared state one wrapper process feeds; every
// access takes the mutex, reads (hits) included.
type responseCache struct {
mu sync.Mutex
cache map[string][]weft.ModelEvent
hits int
}
func (c *responseCache) Hits() int {
c.mu.Lock()
defer c.mu.Unlock()
return c.hits
}
type cacheModel struct {
next weft.Model
c *responseCache
key func(weft.ModelRequest) string
}
func (m cacheModel) Info() weft.ModelInfo { return weft.InfoOf(m.next) }
func (m cacheModel) Stream(ctx context.Context, req weft.ModelRequest) iter.Seq2[weft.ModelEvent, error] {
k := m.key(req)
m.c.mu.Lock()
if events, ok := m.c.cache[k]; ok {
m.c.hits++
m.c.mu.Unlock()
return func(yield func(weft.ModelEvent, error) bool) {
for _, ev := range events {
if !yield(ev, nil) {
return
}
}
}
}
m.c.mu.Unlock()
return func(yield func(weft.ModelEvent, error) bool) {
var seen []weft.ModelEvent
for ev, err := range m.next.Stream(ctx, req) {
if err != nil {
yield(nil, err)
return
}
seen = append(seen, ev)
if !yield(ev, nil) {
return
}
}
m.c.mu.Lock()
m.c.cache[k] = seen
m.c.mu.Unlock()
}
}
func toolNamesOf(tools []*weft.ToolDef) []string {
out := make([]string, len(tools))
for i, t := range tools {
out[i] = t.Name
}
return out
}
Output: provider calls made: 1 of 3
Example (TokenLimitMiddleware) ¶
A token limiter is model middleware: it refuses the call before the provider bills it when the transcript the request carries is over budget (bytes as a proxy — the real count arrives with the finish, when UsageLimit already covers it).
package main
import (
"context"
"fmt"
"iter"
"strings"
"github.com/weftgo/weft"
"github.com/weftgo/weft/wefttest"
)
func main() {
const maxBytes = 64
limiter := func(next weft.Model) weft.Model {
return limitedModel{next: next, max: maxBytes}
}
agt := weft.New(wefttest.Script(wefttest.Say("done")),
weft.WrapModel(limiter))
if _, err := agt.Generate(context.Background(), weft.Prompt(strings.Repeat("x ", 100))); err == nil {
fmt.Println("the oversized prompt went through")
} else {
fmt.Println("refused:", strings.TrimPrefix(err.Error(), "weft: run failed at step 0: model stream: "))
}
}
type limitedModel struct {
next weft.Model
max int
}
func (m limitedModel) Info() weft.ModelInfo { return weft.InfoOf(m.next) }
func (m limitedModel) Stream(ctx context.Context, req weft.ModelRequest) iter.Seq2[weft.ModelEvent, error] {
var n int
for _, msg := range req.Messages {
n += len(msg.Text())
}
if n > m.max {
return func(yield func(weft.ModelEvent, error) bool) {
yield(nil, fmt.Errorf("token limit: request transcript is %d bytes, over the %d-byte budget", n, m.max))
}
}
return m.next.Stream(ctx, req)
}
Output: refused: token limit: request transcript is 200 bytes, over the 64-byte budget
Index ¶
- Variables
- func Allow(permit func(weft.Call) bool) weft.ToolMiddleware
- func Audit(l *slog.Logger) weft.ToolMiddleware
- func Fallback(models ...weft.Model) weft.ModelMiddleware
- func FallbackWhen(when func(error) bool, models ...weft.Model) weft.ModelMiddleware
- func HTTPStatus(err error) (int, bool)
- func Log(l *slog.Logger) weft.ModelMiddleware
- func MapErrors(fn func(error) error) weft.ToolMiddleware
- func RepairJSON() weft.ModelMiddleware
- func Retry(opts ...RetryOption) weft.ModelMiddleware
- func RetryAfter(err error, now time.Time) (time.Duration, bool)
- func Retryable(err error) bool
- type RetryOption
Examples ¶
Constants ¶
This section is empty.
Variables ¶
var ErrRetryAfterTooLong = errors.New("mw: provider asked to retry after longer than the maximum wait")
ErrRetryAfterTooLong is wrapped around a provider error whose retry-after ask exceeds MaxWait: rather than sleeping silently for minutes, Retry fails fast so an outer layer (a queue, a person) can decide. errors.Is on it, or on the provider's own error underneath.
Functions ¶
func Allow ¶
func Allow(permit func(weft.Call) bool) weft.ToolMiddleware
Allow gates every tool call on a predicate over its weft.Call. A denied call never runs: the model sees the error result `DENIED: tool "x" is not allowed` and the loop moves on — there is no retry, no prompt, no exception. Approval is a different seam (weft.RequireApproval): Allow is the fixed policy, approval the deferred decision. Outside the loop (Agent.CallTool with a bare ctx) the Call is built from the call part alone. A nil permit allows every call — the middleware is a pass-through, so a policy that is configured off composes without a branch at the call site.
Example ¶
Allow is a fixed policy: a denied call never runs and the model sees why.
package main
import (
"context"
"fmt"
"log"
"github.com/weftgo/weft"
"github.com/weftgo/weft/mw"
"github.com/weftgo/weft/wefttest"
)
func main() {
rm := weft.Tool("rm", "", func(_ context.Context, _ struct{}) (string, error) { return "gone", nil })
agt := weft.New(wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "rm"}),
wefttest.Say("I could not remove it."),
), rm, weft.WrapTools(mw.Allow(func(c weft.Call) bool { return c.Name != "rm" })))
res, err := agt.Generate(context.Background(), weft.Prompt("remove it"))
if err != nil {
log.Fatal(err)
}
fmt.Println(res.Steps[0].Results[0].Content)
}
Output: DENIED: tool "rm" is not allowed
func Audit ¶
func Audit(l *slog.Logger) weft.ToolMiddleware
Audit logs one line per tool call at Info level on the given logger (slog.Default when nil): run, step, call id, tool, duration, and the outcome — the model-visible error text, plus the internal cause when the error is a *weft.ToolError with one. It changes nothing.
Example ¶
Audit writes one Info line per tool call: run, step, call, tool, duration, and the outcome — with the internal cause when the error is a coded ToolError carrying one.
package main
import (
"bytes"
"context"
"errors"
"fmt"
"log"
"log/slog"
"strings"
"github.com/weftgo/weft"
"github.com/weftgo/weft/mw"
"github.com/weftgo/weft/wefttest"
)
func main() {
var buf bytes.Buffer
logger := slog.New(slog.NewTextHandler(&buf, &slog.HandlerOptions{
ReplaceAttr: func(_ []string, a slog.Attr) slog.Attr {
switch a.Key {
case slog.TimeKey, "dur":
return slog.Attr{}
}
return a
},
}))
charge := weft.Tool("charge", "", func(_ context.Context, _ struct{}) (string, error) {
return "", weft.Errorf("CARD_DECLINED", "the card was declined: %w", errors.New("issuer said no"))
})
agt := weft.New(wefttest.Script(
wefttest.ToolCalls(wefttest.Call{ID: "c1", Name: "charge"}),
wefttest.Say("I could not charge the card."),
), charge, weft.WrapTools(mw.Audit(logger)))
if _, err := agt.Generate(context.Background(), weft.RunID("run-7"), weft.Prompt("charge it")); err != nil {
log.Fatal(err)
}
for _, line := range strings.Split(strings.TrimSpace(buf.String()), "\n") {
fmt.Println(line[strings.Index(line, "msg="):])
}
}
Output: msg="tool call" run=run-7 step=0 call=c1 tool=charge err="CARD_DECLINED: the card was declined: issuer said no" cause="issuer said no"
func Fallback ¶
func Fallback(models ...weft.Model) weft.ModelMiddleware
Fallback tries the given models, in order, when the wrapped model's stream fails before yielding any event — a provider outage, a capability gap (weft.ErrUnsupported). A failure after events were yielded is not retried anywhere: the loop already consumed part of the reply, so it surfaces as the run error. Context cancellation and the WEFT_MODEL_REQUESTS kill switch never fall through. A max_tokens finish is a successful stream, not a failure — Fallback does not switch on it. Info reports the primary model.
Example ¶
Fallback tries another model when the primary fails before yielding.
package main
import (
"context"
"fmt"
"log"
"github.com/weftgo/weft"
"github.com/weftgo/weft/mw"
"github.com/weftgo/weft/wefttest"
)
func main() {
primary := wefttest.Script(wefttest.Fail(fmt.Errorf("%w: pdf input", weft.ErrUnsupported)))
backup := wefttest.Script(wefttest.Say("from the backup model"))
agt := weft.New(primary, weft.WrapModel(mw.Fallback(backup)))
res, err := agt.Generate(context.Background(), weft.Prompt("hi"))
if err != nil {
log.Fatal(err)
}
fmt.Println(res.Text())
}
Output: from the backup model
func FallbackWhen ¶
FallbackWhen is Fallback with a predicate: the next model is tried only for errors that satisfy when. A nil predicate uses Fallback's default (anything but cancellation and the kill switch).
func HTTPStatus ¶
HTTPStatus finds an HTTP status code on the error chain: the vendor SDKs' error types carry it as an exported StatusCode (openai-go, anthropic-sdk-go) or Code (google genai) integer field. The lookup is structural — mw imports no vendor SDK — and reports ok=false when no error on the chain has such a field.
func Log ¶
func Log(l *slog.Logger) weft.ModelMiddleware
Log records one line per model call at Debug level — the request summary before the call and the finish (or error) after it — on the given logger, or slog.Default when nil. Lines carry the model, the message and tool counts, and on finish the stop reason, usage, tool call count, and duration. Placed outermost it sees the outcome of retries and fallbacks; placed innermost, each attempt.
Example ¶
Log writes one Debug line per model call — the request before, the finish (or error) after. Placed outermost it sees the outcome of retries and fallbacks; placed innermost, each attempt.
package main
import (
"bytes"
"context"
"fmt"
"log"
"log/slog"
"github.com/weftgo/weft"
"github.com/weftgo/weft/mw"
"github.com/weftgo/weft/wefttest"
)
func main() {
var buf bytes.Buffer
logger := slog.New(slog.NewTextHandler(&buf, &slog.HandlerOptions{
Level: slog.LevelDebug,
// Drop the volatile attributes so the example's output is stable.
ReplaceAttr: func(_ []string, a slog.Attr) slog.Attr {
switch a.Key {
case slog.TimeKey, "dur":
return slog.Attr{}
}
return a
},
}))
echo := weft.Tool("echo", "", func(_ context.Context, _ struct{}) (string, error) { return "e", nil })
_, err := weft.New(wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "echo"}),
wefttest.Say("ok"),
), echo, weft.WrapModel(mw.Log(logger))).Generate(context.Background(), weft.Prompt("hi"))
if err != nil {
log.Fatal(err)
}
fmt.Println(buf.String())
}
Output: level=DEBUG msg="model request" provider=wefttest model=script messages=1 tools=1 system_bytes=0 thinking=0 level=DEBUG msg="model finish" provider=wefttest model=script reason=tool_calls raw="" input_tokens=10 output_tokens=5 tool_calls=1 level=DEBUG msg="model request" provider=wefttest model=script messages=3 tools=1 system_bytes=0 thinking=0 level=DEBUG msg="model finish" provider=wefttest model=script reason=stop raw="" input_tokens=10 output_tokens=5 tool_calls=0
func MapErrors ¶
func MapErrors(fn func(error) error) weft.ToolMiddleware
MapErrors shapes handler errors centrally before the loop renders them for the model. Errors that are already *weft.ToolError pass through untouched, as do context errors and ErrApprovalRequired (the loop reads those). Every other error goes through fn; a nil fn installs the default, which hides the text — database messages, stack fragments, wrapped internals — behind `INTERNAL: tool "x" failed` while keeping the original as the ToolError's cause for Audit and errors.As.
Example ¶
MapErrors codes handler errors centrally; the default hides the text of unstructured errors behind INTERNAL while keeping the cause for Audit.
package main
import (
"context"
"errors"
"fmt"
"log"
"github.com/weftgo/weft"
"github.com/weftgo/weft/mw"
"github.com/weftgo/weft/wefttest"
)
func main() {
db := weft.Tool("query", "", func(_ context.Context, _ struct{}) (string, error) {
return "", errors.New("pq: SSL is not enabled on the server")
})
agt := weft.New(wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "query"}),
wefttest.Say("The database is unavailable."),
), db, weft.WrapTools(mw.MapErrors(nil)))
res, err := agt.Generate(context.Background(), weft.Prompt("how many orders?"))
if err != nil {
log.Fatal(err)
}
fmt.Println(res.Steps[0].Results[0].Content)
}
Output: INTERNAL: tool "query" failed
func RepairJSON ¶
func RepairJSON() weft.ModelMiddleware
RepairJSON re-encodes a tool call's arguments once when they are not valid JSON — the common damage is a reply cut mid-object or wrapped in a Markdown fence. Unterminated strings, objects, and arrays are closed and fences stripped; when the result parses, it replaces the original. Arguments that cannot be repaired pass through unchanged and fail decoding as they would have, so the model still sees an INVALID_INPUT result. Only tool-call events are touched.
Example ¶
RepairJSON closes a tool call's arguments when the model sends them cut or fenced — the repair the model cannot do for itself. Arguments that cannot be repaired pass through and fail decoding as usual.
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"github.com/weftgo/weft"
"github.com/weftgo/weft/mw"
"github.com/weftgo/weft/wefttest"
)
func main() {
model := wefttest.Script(
wefttest.ToolCalls(wefttest.Call{Name: "city", Args: `{"city":"Par`}),
wefttest.Say("ok"),
)
city := weft.RawTool("city", "Look up a city.", nil, func(_ context.Context, args json.RawMessage) (string, error) {
return "the args arrived as " + string(args), nil
})
res, err := weft.New(model, city, weft.WrapModel(mw.RepairJSON())).Generate(context.Background(), weft.Prompt("x"))
if err != nil {
log.Fatal(err)
}
fmt.Println(res.Steps[0].Results[0].Content)
}
Output: the args arrived as {"city":"Par"}
func Retry ¶
func Retry(opts ...RetryOption) weft.ModelMiddleware
Retry retries a model call that fails before yielding any event — the request-level failures a provider reports up front: rate limits, overload, transport errors, an idle stream. The vendor SDKs already retry at the transport layer (their MaxRetries adapter options); Retry is the loop-visible logic layer above them, with backoff the caller can see in Log and a retry-after the provider sends honoured (retry-after-ms, retry-after, and x-should-retry headers). A failure after events were yielded is never retried: part of the reply has already reached the loop. Exhausted retries return the last error wrapped, so errors.As on the SDK's type still works. Pair Retry with the adapter's MaxRetries(0): the SDK's retries sleep on the provider's retry-after inside the adapter's idle timer (the wait for response headers is one gap), so a long ask fails ErrStreamIdle with the 429 discarded — here the same ask is honoured, capped by MaxWait, and visible in Log.
Example ¶
Retry sits above the vendor SDK's transport retries: it retries the whole model call on request-level failures the classifier accepts.
package main
import (
"context"
"errors"
"fmt"
"log"
"github.com/weftgo/weft"
"github.com/weftgo/weft/mw"
"github.com/weftgo/weft/wefttest"
)
func main() {
model := wefttest.Script(
wefttest.Fail(errors.New("503: overloaded")),
wefttest.Say("hello"),
)
agt := weft.New(model, weft.WrapModel(
mw.Retry(
mw.MaxRetries(2),
mw.BaseDelay(0), // tests only; the default is 500ms doubling to 8s
mw.Classifier(func(err error) bool { return err.Error() == "503: overloaded" }),
),
))
res, err := agt.Generate(context.Background(), weft.Prompt("hi"))
if err != nil {
log.Fatal(err)
}
fmt.Println(res.Text(), len(model.Requests()))
}
Output: hello 2
func RetryAfter ¶
RetryAfter extracts the provider's retry-after ask from the error's HTTP response headers: retry-after-ms (milliseconds), then retry-after (seconds, or an HTTP date relative to now). ok is false when the error carries no such header, or carries one whose value cannot become a duration — unparseable, negative, or too large to convert (ParseFloat accepts 1e19 and Infinity; the int64 conversion of those wraps negative, which would slip past the maxWait cap and sleep zero, turning a misbehaving gateway's ask into up to MaxRetries+1 back-to-back requests).
func Retryable ¶
Retryable is Retry's default classifier: true for weft.ErrStreamIdle, net.Error values, and HTTP 408, 409, 429 and 5xx statuses found on the error (see HTTPStatus); false for context errors, the kill switch, weft.ErrUnsupported, every other 4xx, and errors whose text names a context-window overflow. An x-should-retry header, when the provider sends one, overrides the status rule.
Types ¶
type RetryOption ¶
type RetryOption func(*retryConfig)
RetryOption configures Retry.
func BaseDelay ¶
func BaseDelay(d time.Duration) RetryOption
BaseDelay sets the first backoff delay (default 500ms). Delays double per attempt — 0.5s, 1s, 2s, … — capped at 8s, each with ±25% jitter, unless the provider names its own retry-after, which wins. Zero is meaningful (no delay between attempts); negative values are ignored.
func Classifier ¶
func Classifier(fn func(error) bool) RetryOption
Classifier replaces the default retryability test. The default retries weft.ErrStreamIdle, net.Error values, and HTTP 408, 409, 429 and 5xx responses reported by the vendor SDKs' error types; it never retries context cancellation, the kill switch, weft.ErrUnsupported, quota and billing 4xx, or a context-window overflow — the last must route to compaction, not to the same request again.
func MaxRetries ¶
func MaxRetries(n int) RetryOption
MaxRetries sets how many times a failed call is retried (default 3). Zero disables retrying while keeping the retry-after and classifier behaviour observable through Log. Negative values are ignored — the core options' rule, stated here too.
func MaxWait ¶
func MaxWait(d time.Duration) RetryOption
MaxWait bounds a provider's retry-after ask (default 60s): a longer ask fails fast wrapping ErrRetryAfterTooLong instead of sleeping. MaxWait(0) removes the cap — whatever the provider asks, Retry waits — and negative values are ignored.