go kit
共享基础设施 走吧-\* MCP服务器。一个模块,零膨胀。
go get github.com/anatolykoptev/go-kit包裹
| 包装 | 内容 | 描述 |
|---|---|---|
env | 环境变量解析 | stdlib |
llm | OpenAI兼容的LLM客户端,具有流式传输、工具调用和结构化输出 | stdlib |
cache | L1内存+L2 Redis分层缓存,具有S3-FIFO驱逐、标签无效、字节有界(称重器+最大重量)和空闲(IdleTL)驱逐 | stdlib(L2:Redis) |
retry | 使用指数回退的通用重试 | stdlib |
metrics | 原子计数器、仪表、计时器、标签、接收器、速率、直方图、TTL | stdlib |
hedge | 对冲请求——竞争主请求与备份请求,第一个成功者获胜 | stdlib |
ratelimit | 令牌桶速率限制器、每键支持、并发限制器 | stdlib |
strutil | 支持Unicode的字符串助手,支持大小写转换 | stdlib |
fileopt | 通过gs+qpdf/oxipng/cwebp子流程包装器进行无损PDF/PNG/WebP字节级优化,每个阶段使用Prometheus指标 | stdlib+Prometheus/client_golang |
breaker | 具有指数冷却、抖动、探针槽的3态断路器, Execute[T] 通用包装器, HTTPDoer 预设,按键 Pool | stdlib |
eventbus | 带有点分隔主题和通配符模式匹配的进程中发布/订阅(*, **);64个插槽缓冲通道,直接使用完整语义 | stdlib |
rerank | Cohere兼容的交叉编码器为嵌入式服务器/TEI/Cohere/Jina/Voyage/Mixedbread重新排序HTTP客户端。尽力而为——任何错误都会使输入保持不变。 | stdlib+prometheus/client_golang |
sparse | 用于学习稀疏嵌入的SPLADE形状HTTP客户端——中间件堆栈镜像嵌入/重库(缓存、电路、重试、挂钩、回退)。与密集嵌入/混合检索配对。 | stdlib+prometheus/client_golang |
所有包都是独立的,没有内部交叉导入。只进口你需要的东西。
______________________________________________________________________
环境
import "github.com/anatolykoptev/go-kit/env"
port := env.Int("PORT", 8080)
debug := env.Bool("DEBUG", false)
// Docker secrets / Kubernetes volumes
dbPass := env.File("DB_PASSWORD_FILE", "")
// Variable expansion
dbURL := env.Expand("DATABASE_URL", "postgres://localhost:5432/mydb")
// Binary data
cert := env.Base64("TLS_CERT", nil)
key := env.Hex("API_KEY_HEX", nil)
// Testability — decouple from os.Getenv
env.DefaultSource = env.MapSource(map[string]string{
"PORT": "9090",
})功能: Str, Int, Int64, Uint, Uint64, Float, Bool, Duration, List, Int64List, Map, URL, File, Expand, Base64, Hex.
- 可测试性源接口(用于并行安全测试的MapSource)
- 文件:读取Docker机密和Kubernetes卷
- 展开:解析${VAR}引用
- Base64/Hex:来自环境变量的二进制数据
headers := env.Map("HEADERS", "") // "Content-Type:json,Accept:*/*" → map
endpoint := env.URL("API_URL", "http://localhost:8080") // parsed *url.URL
maxConns := env.Uint("MAX_CONNS", 100)错误处理
// Error-returning variants — return ParseError on invalid values
port, err := env.IntE("PORT", 8080) // err if PORT="abc"
debug, err := env.BoolE("DEBUG", false) // err if DEBUG="maybe"
timeout, err := env.DurationE("TIMEOUT", 30*time.Second) // accepts "5s", "100ms", "2m30s"
// Required — must be set, returns NotSetError if missing
dbURL, err := env.Required("DATABASE_URL")
// Lookup — distinguish "not set" from "set to empty"
val, ok := env.Lookup("OPTIONAL_VAR")
// Must* — panic on invalid (for fail-fast main() init)
dbURL := env.MustRequired("DATABASE_URL")
port := env.MustInt("PORT", 8080)大语言模型
import "github.com/anatolykoptev/go-kit/llm"
client := llm.NewClient(baseURL, apiKey, model,
llm.WithFallbackKeys([]string{backupKey}),
llm.WithMaxTokens(16384),
llm.WithTemperature(0.1),
)
// Simple text completion (unchanged)
response, err := client.Complete(ctx, systemPrompt, userPrompt)
// Full chat with tool calling
resp, err := client.Chat(ctx, messages,
llm.WithTools([]llm.Tool{
llm.NewTool("get_weather", "Get weather for a city", params),
}),
)
for _, call := range resp.ToolCalls { ... }
fmt.Printf("Tokens: %d\n", resp.Usage.TotalTokens)
// Structured output — auto-generates JSON Schema from struct
// Schema constraint tags enrich the JSON Schema for better LLM output:
type User struct {
Name string `json:"name" jsonschema:"description=Full legal name"`
Age int `json:"age" jsonschema:"minimum=0,maximum=150"`
Role string `json:"role" jsonschema:"enum=admin|user|guest"`
}
var recipe Recipe
err := client.ChatTyped(ctx, messages, &recipe)
// SSE streaming
stream, err := client.Stream(ctx, messages)
defer stream.Close()
for chunk, ok := stream.Next(); ok; chunk, ok = stream.Next() {
fmt.Print(chunk.Delta)
}
// Structured extraction with validation retry (Instructor-style)
type Recipe struct {
Name string `json:"name"`
Ingredients []string `json:"ingredients"`
}
var recipe Recipe
err := client.Extract(ctx, messages, &recipe,
llm.WithValidator(func(v any) error {
r := v.(*Recipe)
if len(r.Ingredients) == 0 {
return errors.New("need at least one ingredient")
}
return nil
}),
)
// Union types — LLM chooses between multiple response types
type SearchAction struct {
Query string `json:"query"`
}
type AnswerAction struct {
Answer string `json:"answer"`
}
result, err := client.ExtractOneOf(ctx, messages, []llm.VariantDef{
llm.Variant("search", SearchAction{}),
llm.Variant("answer", AnswerAction{}),
})
switch v := result.(type) {
case *SearchAction:
fmt.Println("Search:", v.Query)
case *AnswerAction:
fmt.Println("Answer:", v.Answer)
}
// Model-level fallback chains
client = llm.NewClient("", "", "",
llm.WithEndpoints([]llm.Endpoint{
{URL: geminiURL, Key: key1, Model: "gemini-2.5-flash"},
{URL: openaiURL, Key: key2, Model: "gpt-4o"},
}),
)
// Request/response middleware
client = llm.NewClient(baseURL, apiKey, model,
llm.WithMiddleware(func(ctx context.Context, req *llm.ChatRequest,
next func(context.Context, *llm.ChatRequest) (*llm.ChatResponse, error)) (*llm.ChatResponse, error) {
start := time.Now()
resp, err := next(ctx, req)
log.Printf("LLM call took %v", time.Since(start))
return resp, err
}),
)- 结构性错误:
APIError{StatusCode, Type, Body, Retryable}--使用errors.As根据错误类型进行分支 - 在429/5xx上以指数回退重试
- 自动回退键循环
- SSE流媒体
Stream/Next - 通过调用工具/功能
Chat+WithTools - 结构化输出通过
ChatTyped+自动JSON模式 - 使用验证重试进行提取(讲师风格,没有Go库这样做)
- 工会类型通过
ExtractOneOf--LLM在响应变量之间进行选择 - 模型级端点回退链
- 用于日志记录、指标、缓存的请求/响应中间件
- 令牌使用情况报告
ChatResponse - 通过以下方式提供多模式支持
CompleteMultimodal - 通过以下方式从LLM输出中提取JSON
ExtractJSON - 架构约束标记:
jsonschema:"description=...,minimum=0,enum=a|b|c"更丰富的模式
缓存
import "github.com/anatolykoptev/go-kit/cache"
// L1-only (no Redis dependency at runtime)
c := cache.New(cache.Config{
L1MaxItems: 1000,
L1TTL: 30 * time.Minute,
})
// L1 + L2 Redis (read-through, write-through)
c := cache.New(cache.Config{
L1MaxItems: 1000,
L1TTL: 30 * time.Minute,
RedisURL: "redis://localhost:6379",
RedisDB: 0,
Prefix: "myapp:",
L2TTL: 24 * time.Hour,
})
// Custom L2 store (testing or alternative backends)
c := cache.New(cache.Config{L1MaxItems: 100, L1TTL: time.Minute})
c.WithL2(myCustomStore)
defer c.Close()
c.Set(ctx, "key", data)
data, ok := c.Get(ctx, "key")
// Cache-aside with singleflight (concurrent loads deduplicated)
data, err := c.GetOrLoad(ctx, "key", func(ctx context.Context) ([]byte, error) {
return fetchFromDB(ctx, "key")
})
// Statistics
stats := c.Stats()
fmt.Printf("Hit ratio: %.1f%%, Evictions: %d\n", stats.HitRatio*100, stats.Evictions)按键TTL --为单个条目覆盖全局TTL:
// Short TTL for fast-changing data (e.g. job listings)
c.SetWithTTL(ctx, "jobs:123", data, 15*time.Minute)
// Cache-aside with custom TTL
data, err := c.GetOrLoadWithTTL(ctx, "company:456", 24*time.Hour,
func(ctx context.Context) ([]byte, error) {
return fetchCompanyData(ctx, "456")
},
)- 具有S3-FIFO逐出功能的L1内存缓存,可实现高命中率
- L2 Redis:可选,优雅降级(仅当Redis无法访问时为L1)
- 通读:L1错过→ L2命中→ 自动L1升级
- 直写:设置/删除会传播到两个层
- L2接口:插入自定义后端进行测试或替代
- GetOrLoad带直列式单飞行(防止成群打雷)
- TTL抖动(防止缓存踩踏)
- 统计数据中的驱逐计数器+点击率
- 后台清理,TTL到期
- 驱逐通知的OnEvict回调(过期、容量、显式)
- 基于标签的失效:按标签对条目进行分组,批量失效
- 类型化JSON助手:泛型
SetJSON/GetJSON/GetOrLoadJSON
基于标签的失效 --对相关条目进行分组并使其无效:
c.SetWithTags(ctx, "user:1:profile", data, []string{"user:1", "profile"})
c.SetWithTags(ctx, "user:1:settings", data, []string{"user:1"})
n := c.InvalidateByTag(ctx, "user:1") // removes both entries, returns 2
tags := c.Tags("user:1:profile") // []string{"user:1", "profile"}类型化JSON缓存 --通用包装 []byte API
cache.SetJSON(c, ctx, "user:1", User{Name: "Alice", Age: 30})
user, ok, err := cache.GetJSON[User](c, ctx, "user:1")
user, err := cache.GetOrLoadJSON[User](c, ctx, "user:1", func(ctx context.Context) (User, error) {
return fetchUser(ctx, 1)
})OnEvict回拨 --对缓存驱逐做出反应:
c := cache.New(cache.Config{
L1MaxItems: 1000,
L1TTL: 30 * time.Minute,
OnEvict: func(key string, data []byte, reason cache.EvictReason) {
switch reason {
case cache.EvictCapacity:
metrics.Incr("cache.evict.capacity")
case cache.EvictExpired:
metrics.Incr("cache.evict.expired")
case cache.EvictExplicit:
metrics.Incr("cache.evict.explicit")
}
},
})对冲
import "github.com/anatolykoptev/go-kit/hedge"
// Start fn; if no response after 1s, launch a second call in parallel.
// First success wins, loser is cancelled automatically.
result, err := hedge.Do(ctx, time.Second, func(ctx context.Context) (string, error) {
return callLLM(ctx)
})
// Zero/negative delay: run fn once, no goroutines.
result, err := hedge.Do(ctx, 0, fn)- 通用的
Do[T any]--适用于任何返回类型 - 共享派生上下文--
defer cancel()自动清理失败者goroutine - 主延迟前失败——立即返回错误,无对冲
- 缓冲通道可防止goroutine泄漏
速率限制
import "github.com/anatolykoptev/go-kit/ratelimit"
// Single rate limiter: 10 requests/sec, burst of 5
lim := ratelimit.New(10, 5)
if lim.Allow() {
// proceed
}
// Blocking wait (respects context cancellation)
err := lim.Wait(ctx)
// Per-key rate limiting (per-domain, per-API-key)
kl := ratelimit.NewKeyLimiter(5, 3) // 5/sec per key, burst 3
defer kl.Close()
kl.Allow("api.linkedin.com")
kl.Wait(ctx, "api.twitter.com")
// Background cleanup of idle limiters
kl.StartCleanup(time.Minute, 10*time.Minute)并发限制器 (舱壁样式):
// Limit to 5 concurrent operations
cl := ratelimit.NewConcurrencyLimiter(5)
release, err := cl.Acquire(ctx) // blocking; respects context
if err != nil { return err }
defer release()
// Non-blocking variant
release, ok := cl.TryAcquire()
cl.Available() // free slots
cl.Size() // max slots- 令牌桶算法,零外部deps
- 非阻塞
Allow()和封锁Wait(ctx) - 具有自动空闲清理功能的按键限制器
- 并发限制器(基于信号量,阻塞+非阻塞获取)
- Goroutine保险箱
重试
import "github.com/anatolykoptev/go-kit/retry"
result, err := retry.Do(ctx, retry.Options{
MaxAttempts: 5,
InitialDelay: 500 * time.Millisecond,
MaxDelay: 10 * time.Second,
MaxElapsedTime: 30 * time.Second, // total budget
Jitter: true, // ±25% random jitter
}, func() (string, error) {
return callAPI()
})
// HTTP-specific: retries on 429/5xx, auto-parses Retry-After header
resp, err := retry.HTTP(ctx, retry.Options{Jitter: true}, doRequest)
// Override backoff from fn:
return "", retry.RetryAfter(5*time.Second, err)
// Abort on specific errors (never retry)
retry.Do(ctx, retry.Options{
AbortOn: []error{context.DeadlineExceeded},
}, fn)
// Opt-in retry: only marked errors are retried
retry.Do(ctx, retry.Options{RetryableOnly: true}, func() (T, error) {
return result, retry.MarkRetryable(err) // will retry
})
// Permanent error — stop retrying immediately
retry.Do(ctx, retry.Options{MaxAttempts: 5}, func() (T, error) {
if isFatal(err) {
return zero, retry.Permanent(err) // unwrapped and returned
}
return zero, err
})
// OnRetry callback — log each failed attempt
retry.Do(ctx, retry.Options{
MaxAttempts: 5,
OnRetry: func(attempt int, err error) {
log.Printf("attempt %d failed: %v", attempt, err)
},
}, fn)
// RetryIf — custom predicate (overrides AbortOn + RetryableOnly)
retry.Do(ctx, retry.Options{
MaxAttempts: 5,
RetryIf: func(err error) bool {
var netErr net.Error
return errors.As(err, &netErr) && netErr.Temporary()
},
}, fn)- AbortOn:从不重试特定错误(例如context.DaillineExceded)
- 仅可重试+标记可重试:选择重试模式以确保生产安全
- RetryInf:自定义谓词--完全控制要重试的错误
- 永久(err):fn发出立即停止重试的信号
- OnRetry回调:每次失败尝试的日志记录/指标
- 上下文错误包装:
errors.Is(err, context.DeadlineExceeded)超时工作
指标
import "github.com/anatolykoptev/go-kit/metrics"
reg := metrics.NewRegistry()
// Counters
reg.Incr("requests")
reg.Add("bytes", 1024)
// Gauges — track current values
reg.Gauge("connections").Inc()
reg.Gauge("cpu").Set(45.2)
reg.Gauge("queue").Dec()
// Timer — one-liner duration tracking
defer reg.StartTimer("api.latency").Stop()
// Labels — dimensional metrics
reg.Incr(metrics.Label("requests", "method", "GET"))
reg.Incr(metrics.Label("requests", "method", "POST"))
// Rate tracking (EWMA)
rate := reg.Rate("events")
rate.Update(1) // record event
rate.M1() // events/sec, 1-minute window
// Histogram (percentiles via reservoir sampling)
h := reg.Histogram("latency")
h.Update(12.5) // record observation
snap := h.Snapshot()
// snap.P50, snap.P95, snap.P99, snap.Min, snap.Max, snap.Mean
// TTL for dynamic metrics
reg.IncrWithTTL(metrics.Label("api.calls", "path", "/users"), 10*time.Minute)
reg.CleanupExpired() // remove stale metrics
// Snapshot and reset (for periodic reporting)
all := reg.SnapshotAndReset() // reads + zeros atomically
// Output formatting
reg.WriteTo(os.Stdout, metrics.TextSink{}) // key=value lines
reg.WriteTo(os.Stdout, metrics.JSONSink{}) // JSON object- 带无锁定浮动的仪表类型64(设置/添加/增加/减少)
- StartTimer/Stop用于单线延迟跟踪
- 维度度量键的标签()
- 速率(EWMA):事件/秒,1/5/15分钟移动平均线
- 柱状图:P50/P95/P99的储层采样,无无限内存
- TTL:自动过期的每端点/每用户指标
- SnapshotAndReset用于原子读取和归零
- 带有TextSink和JSONSink的接收器接口
结构
import "github.com/anatolykoptev/go-kit/strutil"
s := strutil.Truncate("Hello, world!", 5) // "Hello..."
s = strutil.TruncateAtWord("Hello, world!", 8) // "Hello,..."
s = strutil.TruncateMiddle("path/to/file.go", 10) // "path/...e.go"
// Custom placeholder
s = strutil.TruncateWith("Hello, world!", 5, "[...]") // "Hello[...]"
// Case conversions
s = strutil.ToSnakeCase("myVariableName") // "my_variable_name"
s = strutil.ToCamelCase("my_variable") // "myVariable"
s = strutil.ToKebabCase("myVariableName") // "my-variable-name"
s = strutil.ToPascalCase("my_variable") // "MyVariable"
// Word wrap
wrapped := strutil.WordWrap("long text here...", 80)
// Clean invalid UTF-8
clean := strutil.Scrub(untrustedInput)
// Check all substrings present
strutil.ContainsAll(s, []string{"foo", "bar"})
ok := strutil.Contains([]string{"a", "b"}, "a") // true
ok = strutil.ContainsAny("hello world", []string{"world"}) // true- WordWrap:在单词边界处包裹文本
- Scrub:用U+FFFD替换无效的UTF-8
- ContainsAll:检查是否存在所有子字符串
消费者
| 服务 | 使用的包 |
|---|---|
| 去搜索 | 缓存、环境、字符串 |
| 去工作 | 缓存、环境、llm、指标、结构 |
| 去wp | 缓存、环境、llm、指标、结构 |
| go代码 | 缓存、环境、llm |
| 加油 | 缓存、环境、llm、指标、结构 |
| 启动 | 缓存、环境、llm、指标、重试、strutil |
| 鼓起勇气 | 环境、生命周期管理、指标 |
| 卫生文本 | 环境、生命周期管理、指标 |
fileopt
通过子进程包装器对PDF/PNG/WebP进行无损字节级优化 gs+qpdf, oxipng,以及 cwebp。专为生成或接收文档并希望在磁盘写入、上传或LLM输入之前减小有效载荷大小的服务而设计。
import "github.com/anatolykoptev/go-kit/fileopt"
// Dispatch by extension
opt, err := fileopt.OptimizeBytes(ctx, data,
fileopt.KindFromExt(filepath.Ext(filename)),
fileopt.LevelEbook, 80)
// Or call specific optimizer
opt, err := fileopt.OptimizePNG(ctx, data)
opt, err := fileopt.OptimizePDF(ctx, data, fileopt.LevelEbook)
opt, err := fileopt.OptimizeWebP(ctx, data, 80)
// Expose Prometheus metrics
mux.Handle("/metrics", fileopt.MetricsHandler())保证:
- 默认情况下是无损的:当一个阶段会增长文件时,大小救援卫士会返回原始值(cwebp梯度反模式)。
- 内容感知:纯文本PDF跳过gs阶段(10-16倍加速;qpdf独自承担工作)。
- 每个阶段普罗米修斯指标:
gokit_fileopt_{calls_total, duration_seconds, ratio, bytes_before_total, bytes_after_total}标记为stage(gs/qpdf/oxipng/cwebp)和result(成功/跳过/错误)。
系统二进制覆盖: FILEOPT_GS_PATH, FILEOPT_QPDF_PATH, FILEOPT_OXIPNG_PATH, FILEOPT_CWEBP_PATH.缺少二进制文件→ 警告日志+原始字节(永远不会使调用者失败)。
重排序
Cohere-shape HTTP客户端,用于跨编码器重新排序端点。兼容 embed-server 自托管、HuggingFace TEI和Cohere/Jina/Voyage/Mixedbread托管提供商。
import "github.com/anatolykoptev/go-kit/rerank"
c := rerank.New(rerank.Config{
URL: "http://embed-server:8082",
Model: "gte-multi-rerank",
Timeout: 4 * time.Second,
MaxDocs: 20,
}, nil)
scored := c.Rerank(ctx, query, []rerank.Doc{
{ID: "u1", Text: "..."},
{ID: "u2", Text: "..."},
})
// scored sorted by .Score desc; .OrigRank preserved; zero Score for docs the server didn't rank.保证:
- 尽力而为:任何错误(超时、非2xx、解码)都会以不变的方式返回输入
slog.Warn并且从不传播error价值——管道总是向前发展。 - 零值
Config.URL禁用客户端;Rerank返回输入不变,Available()回报false. - 头部/尾部分开
MaxDocs:超出上限的文档在重新分级后按原始顺序保存。 MaxCharsPerDoc(支持rune,UTF-8安全)发送到服务器的每个文档文本的边界——防止O(seq²)长输入会引起交叉编码器的注意。
普罗米修斯: rerank_requests_total{model,status} (计数器), rerank_duration_seconds{model} (直方图,桶0.05..10s)。
认证: Config.APIKey 套 Authorization: Bearer --Cohere/Jina/Voyage/Mixedbread托管所需;为自托管嵌入式服务器/TEI留空。
稀疏
用于学习稀疏嵌入端点的SPLADE形状HTTP客户端。镜子 嵌入/ 和 重新排序/ 中间件堆栈(缓存、电路、重试、挂钩、回退)。专为自托管Rust设计 嵌入式服务器 服务于SPLADE-v3-distilbert的侧车;与任何TEI风格兼容 /embed_sparse 服务器返回Qdrant形状稀疏向量包络。
import "github.com/anatolykoptev/go-kit/sparse"
c, _ := sparse.NewClient("http://embed-server:8082",
sparse.WithModel("splade-v3-distilbert"),
sparse.WithTimeout(30*time.Second),
)
defer c.Close()
vecs, err := c.EmbedSparse(ctx, []string{"first text", "second text"})
// each vec.Indices: BERT vocab token ids (uint32);
// each vec.Values: log(1+ReLU(logit)) weights, sorted by weight desc.为什么稀疏与密集嵌入/并存? SPLADE是BM25+神经术语扩展的一种。Dense(e5,ada)在跨语言语义释义方面获胜;稀疏赢得稀有项(名称、ID、品牌、版本号)并插入倒排索引(pgvector) sparsevec,Qdrant稀疏,Lucene)。混合检索=密集+稀疏+RRF优于单独检索。
服务器合同: POST /embed_sparse --身体 {"input":["..."],"model":"...","top_k":256,"min_weight":0.0},响应 {"model":"...","data":[{"index":N,"indices":[...],"values":[...]}]}. top_k=0 (默认)允许服务器选择其默认值。空输入→ (nil, nil).
弹性(嵌入镜子/):
- 在出现指数回退+抖动的瞬态故障(超时,429,5xx)时重试;4xx上不可重试。
- 可选的Redis L2缓存(默认情况下关闭——稀疏流量由索引主导,其中每个文本只显示一次;通过
WithCache用于查询侧热路径)。 - 可选断路器(
WithCircuit). - 可选初级→回退链(
WithFallback). - 键入错误
ErrModelNotConfigured当传递错误的型号名称时,400路径。
普罗米修斯: gokit_sparse_requests_total{outcome,backend} 计数器, gokit_sparse_request_duration_seconds{backend} 直方图, gokit_sparse_batch_size{backend} 直方图, gokit_sparse_terms_per_vector{backend} 直方图, gokit_sparse_retry_attempt_total{backend,attempt} 柜台,加上标准 gokit_sparse_retry_total{backend,reason}.
环境驱动型建筑(sparse.New(...)): SPARSE_BACKEND=http, SPARSE_HTTP_BASE_URL, SPARSE_MODEL (默认值 splade-v3-distilbert), SPARSE_HTTP_TIMEOUT,可选 SPARSE_TOP_K, SPARSE_MIN_WEIGHT.
许可证
Apache 2.0
