From 6bf027367a500f939decd56ace1e92f438deb157 Mon Sep 17 00:00:00 2001 From: M09Ic Date: Sat, 26 Sep 2026 20:13:00 -0700 Subject: [PATCH] feat(jev): claim contract and Jev client as a standalone module Moves the provider-neutral judgement contract out of fingers/judge so fingers keeps only fingerprint claims: - Claim / Option / Outcome / Ruling, Claim.Resolve, ValidateClaims and ValidateRulings: a Provider picks one finite option per claim with a confidence; Resolve is the only interpretation of that pick. - Provider interface and jev.Judge, which validates a batch, freezes the evidence as JSON once and validates the complete answer. - Cached: exact LRU over whole batches with coalescing of identical concurrent requests; failures are never cached. - Client: the TypeSafe Jev (System One) HTTP provider, with retry on 429/529 and DefaultMinConfidence 0.3 calibrated for jev-1.13.0. go 1.17, standard library only. Co-Authored-By: Claude Opus 5.5 (1M context) --- jev/README.md | 45 ++++++++++++ jev/cache.go | 101 +++++++++++++++++++++++++ jev/cache_test.go | 179 +++++++++++++++++++++++++++++++++++++++++++++ jev/claim.go | 119 ++++++++++++++++++++++++++++++ jev/claim_test.go | 118 ++++++++++++++++++++++++++++++ jev/client.go | 160 ++++++++++++++++++++++++++++++++++++++++ jev/client_test.go | 90 +++++++++++++++++++++++ jev/doc.go | 16 ++++ jev/go.mod | 3 + jev/provider.go | 54 ++++++++++++++ 10 files changed, 885 insertions(+) create mode 100644 jev/README.md create mode 100644 jev/cache.go create mode 100644 jev/cache_test.go create mode 100644 jev/claim.go create mode 100644 jev/claim_test.go create mode 100644 jev/client.go create mode 100644 jev/client_test.go create mode 100644 jev/doc.go create mode 100644 jev/go.mod create mode 100644 jev/provider.go diff --git a/jev/README.md b/jev/README.md new file mode 100644 index 0000000..dd45cc7 --- /dev/null +++ b/jev/README.md @@ -0,0 +1,45 @@ +# jev:以 Claim 向模型求判断 + +代码提出断言(Claim)并附上证据,Provider 从有限选项中选一个并给出置信度,`Claim.Resolve` 把选择解释为 holds / refuted / insufficient。模型不写回自由文本。内置 `Client` 对接 [TypeSafe Jev](https://docs.typesafe.ai/api.md)(System One),其他模型实现 `Provider` 即可接入。 + +```go +c, err := jev.NewClient("") // TYPESAFE_API_KEY +p := jev.Cached(c, jev.DefaultCacheSize) +rulings, err := jev.Judge(ctx, p, state, claims) +outcome := claims["k"].Resolve(rulings["k"], jev.DefaultMinConfidence) +``` + +`jev.Judge` 校验输入,把 state 冻结为 JSON 一次,并校验完整回答;失败的批次不返回任何 Ruling。`DefaultMinConfidence`(0.3)是针对固定模型 `DefaultModel` 在指纹判定验收集上标定的阈值。 + +## 合约 + +```go +type Claim struct { + Statement string + Options map[string]Option +} +type Option struct { + Description string + Outcome Outcome // Insufficient / Holds / Refuted +} +type Ruling struct { + Option string + Confidence float64 +} +type Provider interface { + ID() string + Judge(context.Context, interface{}, map[string]Claim) (map[string]Ruling, error) +} +``` + +每个 Claim 必须有非空 Statement、非空批次 key,以及显式提供的 `insufficient` 选项,映射到 `Insufficient`。每个 Option 都必须声明合法 Outcome,`jev.Judge` 不自动增补或改写。Provider 返回完整且精确对应的 key 集合,选项必须存在,置信度必须是 [0,1] 内有限值。非法或缺失回答是错误,不当作证据不足。`jev.ValidateClaims` 在调用前校验输入,`jev.ValidateRulings` 在应用或缓存前校验完整批次。 + +`Outcome` 直接使用字符串 `holds` / `refuted` / `insufficient`,零值为空表示未判定。Claim 的 JSON 为 `statement/options`,Option 为 `description/outcome`,Ruling 为 `option/confidence`;`questions/answers/choice` 只出现在 Jev HTTP 边界,不进入核心和缓存格式。 + +`claim.Resolve(ruling, minConfidence)` 是唯一的语义解释入口。低置信度解析为 `Insufficient`,但原始 `Ruling.Option` 保留。调用方用代码确定的事实也可以构建合法 Claim/Ruling(置信度 1),经相同的 Resolve 解释。 + +## 缓存与 Provider + +`jev.Cached(p, capacity)` 包装完整批次的精确 LRU,并合并并发相同请求;容量非正时不包装。键包括 Provider.ID、完整提交证据和 Claim 选项及语义,不进行页面相似度复用。失败、非法回答不缓存;返回 map 独立复制。等待者可独立取消,发起者取消会使该次共享请求失败,后续可重试。 + +`jev.Client` 的 ID 包含 endpoint 和 model,推荐以 `jev.Cached(c, jev.DefaultCacheSize)`(4096)包装。用户可以实现自己的 Provider 包装器处理缓存、限速和计数。Client 只将 Statement 和 Option.Description 转为服务端 choice 格式,不发送本地 Outcome 策略;响应必须显式包含 confidence,缺省或 null 会报错。 diff --git a/jev/cache.go b/jev/cache.go new file mode 100644 index 0000000..57ed206 --- /dev/null +++ b/jev/cache.go @@ -0,0 +1,101 @@ +package jev + +import ( + "container/list" + "context" + "crypto/sha256" + "encoding/json" + "sync" +) + +const DefaultCacheSize = 4096 + +// Cached adds an exact, bounded in-memory cache and coalesces identical +// concurrent batches. Wrap Provider yourself for other storage or policies. +// A non-positive capacity disables this wrapper. Failed batches are not cached. +func Cached(p Provider, capacity int) Provider { + if p == nil || capacity <= 0 { + return p + } + return &cachedProvider{Provider: p, capacity: capacity, order: list.New(), + entries: map[[32]byte]*list.Element{}, pending: map[[32]byte]*pendingBatch{}} +} + +type cachedProvider struct { + Provider + capacity int + mu sync.Mutex + order *list.List + entries map[[32]byte]*list.Element + pending map[[32]byte]*pendingBatch +} + +type cacheEntry struct { + key [32]byte + rulings map[string]Ruling +} + +type pendingBatch struct { + done chan struct{} + rulings map[string]Ruling + err error +} + +func (p *cachedProvider) Judge(ctx context.Context, state interface{}, claims map[string]Claim) (map[string]Ruling, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + if err := ValidateClaims(claims); err != nil { + return nil, err + } + data, err := json.Marshal([]interface{}{p.ID(), state, claims}) + if err != nil { + return nil, err + } + key := sha256.Sum256(data) + p.mu.Lock() + if entry := p.entries[key]; entry != nil { + p.order.MoveToFront(entry) + out := copyRulings(entry.Value.(*cacheEntry).rulings) + p.mu.Unlock() + return out, nil + } + if pending := p.pending[key]; pending != nil { + p.mu.Unlock() + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-pending.done: + if pending.err != nil { + return nil, pending.err + } + return copyRulings(pending.rulings), nil + } + } + pending := &pendingBatch{done: make(chan struct{})} + p.pending[key] = pending + p.mu.Unlock() + + got, err := p.Provider.Judge(ctx, state, claims) + if err == nil { + err = ValidateRulings(claims, got) + } + p.mu.Lock() + if err == nil { + pending.rulings = copyRulings(got) + p.entries[key] = p.order.PushFront(&cacheEntry{key, pending.rulings}) + if p.order.Len() > p.capacity { + last := p.order.Back() + delete(p.entries, last.Value.(*cacheEntry).key) + p.order.Remove(last) + } + } + pending.err = err + delete(p.pending, key) + close(pending.done) + p.mu.Unlock() + if err != nil { + return nil, err + } + return copyRulings(pending.rulings), nil +} diff --git a/jev/cache_test.go b/jev/cache_test.go new file mode 100644 index 0000000..86da8ca --- /dev/null +++ b/jev/cache_test.go @@ -0,0 +1,179 @@ +package jev + +import ( + "context" + "errors" + "math" + "sync" + "sync/atomic" + "testing" +) + +type providerFunc struct { + id string + call func(context.Context, interface{}, map[string]Claim) (map[string]Ruling, error) +} + +func (p *providerFunc) ID() string { return p.id } +func (p *providerFunc) Judge(ctx context.Context, s interface{}, c map[string]Claim) (map[string]Ruling, error) { + return p.call(ctx, s, c) +} + +func TestExactCacheIdentityCloneAndEviction(t *testing.T) { + calls := 0 + p := &providerFunc{id: "endpoint/model", call: func(_ context.Context, _ interface{}, cs map[string]Claim) (map[string]Ruling, error) { + calls++ + out := map[string]Ruling{} + for k := range cs { + out[k] = Ruling{Option: "yes", Confidence: 1} + } + return out, nil + }} + cached := Cached(p, 2) + claims := map[string]Claim{"c": testClaim("is it running?")} + ask := func(state string) map[string]Ruling { + t.Helper() + r, e := cached.Judge(context.Background(), state, claims) + if e != nil { + t.Fatal(e) + } + return r + } + first := ask("a") + first["c"] = Ruling{Option: "no"} + if ask("a")["c"].Option != "yes" || calls != 1 { + t.Fatal("cache shares output or misses identical input") + } + ask("b") + ask("a") + ask("c") + ask("b") + if calls != 4 { + t.Fatalf("LRU calls=%d", calls) + } + p.id = "endpoint/other-model" + ask("b") + p.id = "other-endpoint/other-model" + ask("b") + claim := claims["c"] + claim.Options["yes"] = Option{"new evidence meaning", Holds} + claims["c"] = claim + ask("b") + claim.Statement = "another claim" + claims["c"] = claim + ask("b") + if calls != 8 { + t.Fatalf("identity/claim changes reused cached output: %d", calls) + } +} + +func TestCacheDoesNotStoreFailuresOrInvalidBatches(t *testing.T) { + for _, mode := range []string{"error", "missing", "unknown", "nan", "extra"} { + t.Run(mode, func(t *testing.T) { + calls := 0 + p := &providerFunc{call: func(_ context.Context, _ interface{}, cs map[string]Claim) (map[string]Ruling, error) { + calls++ + if calls > 1 { + return map[string]Ruling{"c": {Option: "yes", Confidence: 1}}, nil + } + switch mode { + case "error": + return nil, errors.New("unavailable") + case "missing": + return nil, nil + case "unknown": + return map[string]Ruling{"c": {Option: "invented", Confidence: 1}}, nil + case "nan": + return map[string]Ruling{"c": {Option: "yes", Confidence: math.NaN()}}, nil + default: + return map[string]Ruling{"c": {Option: "yes", Confidence: 1}, "extra": {}}, nil + } + }} + c := Cached(p, 1) + claims := map[string]Claim{"c": testClaim("?")} + if _, err := c.Judge(context.Background(), "state", claims); err == nil { + t.Fatal("invalid batch succeeded") + } + for i := 0; i < 2; i++ { + if _, err := c.Judge(context.Background(), "state", claims); err != nil { + t.Fatal(err) + } + } + if calls != 2 { + t.Fatalf("failed batch cached, or good batch missed: %d", calls) + } + }) + } +} + +func TestCacheConcurrentSharingAndCancellation(t *testing.T) { + started, release := make(chan struct{}), make(chan struct{}) + var calls int64 + p := &providerFunc{call: func(ctx context.Context, _ interface{}, _ map[string]Claim) (map[string]Ruling, error) { + atomic.AddInt64(&calls, 1) + close(started) + select { + case <-release: + return map[string]Ruling{"c": {Option: "yes", Confidence: 1}}, nil + case <-ctx.Done(): + return nil, ctx.Err() + } + }} + c := Cached(p, 2) + claims := map[string]Claim{"c": testClaim("?")} + var wg sync.WaitGroup + wg.Add(1) + go func() { + defer wg.Done() + if _, e := c.Judge(context.Background(), "s", claims); e != nil { + t.Error(e) + } + }() + <-started + ctx, cancel := context.WithCancel(context.Background()) + canceled := make(chan error, 1) + go func() { _, e := c.Judge(ctx, "s", claims); canceled <- e }() + cancel() + if !errors.Is(<-canceled, context.Canceled) { + t.Fatal("waiter did not cancel") + } + for i := 0; i < 20; i++ { + wg.Add(1) + go func() { + defer wg.Done() + r, e := c.Judge(context.Background(), "s", claims) + if e != nil || r["c"].Option != "yes" { + t.Errorf("ruling=%v error=%v", r, e) + } + }() + } + close(release) + wg.Wait() + if calls != 1 { + t.Fatalf("duplicate provider calls: %d", calls) + } + if _, e := c.Judge(ctx, "s", claims); !errors.Is(e, context.Canceled) { + t.Fatal("canceled request got cached success") + } +} + +func TestCanceledLeaderIsNotCached(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + calls := 0 + p := &providerFunc{call: func(ctx context.Context, _ interface{}, _ map[string]Claim) (map[string]Ruling, error) { + calls++ + if calls == 1 { + cancel() + return nil, ctx.Err() + } + return map[string]Ruling{"c": {Option: "yes", Confidence: 1}}, nil + }} + c := Cached(p, 1) + claims := map[string]Claim{"c": testClaim("?")} + if _, e := c.Judge(ctx, "s", claims); !errors.Is(e, context.Canceled) { + t.Fatal(e) + } + if _, e := c.Judge(context.Background(), "s", claims); e != nil || calls != 2 { + t.Fatalf("retry %v %d", e, calls) + } +} diff --git a/jev/claim.go b/jev/claim.go new file mode 100644 index 0000000..ba61ab3 --- /dev/null +++ b/jev/claim.go @@ -0,0 +1,119 @@ +package jev + +import ( + "fmt" + "math" + "strings" +) + +// Claim defines a statement and the meaning of each finite option. Evidence +// is shared by the batch passed to Judge. Every claim explicitly abstains. +type Claim struct { + Statement string `json:"statement"` + Options map[string]Option `json:"options"` +} + +// Option is part of a Claim, not a separate judgement mechanism. +type Option struct { + Description string `json:"description,omitempty"` + Outcome Outcome `json:"outcome"` +} + +// Outcome is a resolved Claim meaning. The zero value means unjudged. +// String values are shared by live results, audit records and JSON snapshots. +type Outcome string + +const ( + Insufficient Outcome = "insufficient" + Holds Outcome = "holds" + Refuted Outcome = "refuted" +) + +func (o Outcome) String() string { return string(o) } +func (o Outcome) valid() bool { return o == Insufficient || o == Holds || o == Refuted } + +// Ruling records the original selection and its confidence. Its outcome is +// always derived from the claim; low confidence does not rewrite the option. +type Ruling struct { + Option string `json:"option"` + Confidence float64 `json:"confidence"` +} + +// OptionInsufficient is the explicit abstention key required by every Claim. +const OptionInsufficient = "insufficient" + +func probability(v float64) bool { + return !math.IsNaN(v) && !math.IsInf(v, 0) && v >= 0 && v <= 1 +} + +func (c Claim) validate() error { + if strings.TrimSpace(c.Statement) == "" { + return fmt.Errorf("a claim requires a statement") + } + if len(c.Options) < 2 { + return fmt.Errorf("a claim requires options and an explicit insufficient option") + } + abstain, ok := c.Options[OptionInsufficient] + if !ok || abstain.Outcome != Insufficient { + return fmt.Errorf("insufficient must explicitly mean Insufficient") + } + for key, option := range c.Options { + if strings.TrimSpace(key) == "" || !option.Outcome.valid() { + return fmt.Errorf("invalid option %q", key) + } + } + return nil +} + +func (c Claim) validateRuling(r Ruling) error { + if _, ok := c.Options[r.Option]; !ok { + return fmt.Errorf("unknown option %q", r.Option) + } + if !probability(r.Confidence) { + return fmt.Errorf("invalid confidence %v", r.Confidence) + } + return nil +} + +// Resolve is the only interpretation of a ruling. Invalid input cannot +// establish a claim; execution reports invalid inputs as errors beforehand. +func (c Claim) Resolve(r Ruling, minConfidence float64) Outcome { + if c.validate() != nil || c.validateRuling(r) != nil || !probability(minConfidence) || r.Confidence < minConfidence { + return Insufficient + } + return c.Options[r.Option].Outcome +} + +// ValidateClaims checks inputs before a Provider is invoked. +func ValidateClaims(claims map[string]Claim) error { + for key, c := range claims { + if strings.TrimSpace(key) == "" { + return fmt.Errorf("jev: empty claim key") + } + if err := c.validate(); err != nil { + return fmt.Errorf("jev: claim %q: %w", key, err) + } + } + return nil +} + +// ValidateRulings checks the complete Claim/Ruling contract, including exact keys. +// Provider wrappers can use it before persisting results. +func ValidateRulings(claims map[string]Claim, rulings map[string]Ruling) error { + if err := ValidateClaims(claims); err != nil { + return err + } + if len(rulings) != len(claims) { + return fmt.Errorf("jev: provider returned %d rulings for %d claims", len(rulings), len(claims)) + } + for key, c := range claims { + r, ok := rulings[key] + if !ok { + return fmt.Errorf("jev: provider omitted claim %q", key) + } + if err := c.validateRuling(r); err != nil { + return fmt.Errorf("jev: claim %q: %w", key, err) + } + } + return nil +} diff --git a/jev/claim_test.go b/jev/claim_test.go new file mode 100644 index 0000000..531be60 --- /dev/null +++ b/jev/claim_test.go @@ -0,0 +1,118 @@ +package jev + +import ( + "context" + "encoding/json" + "math" + "reflect" + "testing" +) + +func TestClaimAndRulingJSON(t *testing.T) { + claim := testClaim("The response is produced by this product.") + data, err := json.Marshal(claim) + if err != nil { + t.Fatal(err) + } + var restored Claim + if err := json.Unmarshal(data, &restored); err != nil || !reflect.DeepEqual(claim, restored) { + t.Fatalf("claim roundtrip: %s, %v", data, err) + } + ruling := Ruling{Option: "yes", Confidence: 0.2} + data, err = json.Marshal(ruling) + if err != nil || string(data) != `{"option":"yes","confidence":0.2}` { + t.Fatalf("provider vocabulary leaked into ruling: %s, %v", data, err) + } + var got Ruling + if err := json.Unmarshal(data, &got); err != nil || got != ruling || restored.Resolve(got, 0.5) != Insufficient || restored.Resolve(got, 0.1) != Holds { + t.Fatalf("ruling or interpretation changed: %+v, %v", got, err) + } + var unjudged Outcome + for _, outcome := range []Outcome{unjudged, Holds, Refuted, Insufficient} { + data, err := json.Marshal(outcome) + var decoded Outcome + if err != nil || json.Unmarshal(data, &decoded) != nil || decoded != outcome { + t.Fatalf("outcome roundtrip: %q, %s, %v", outcome, data, err) + } + } + if unjudged == Insufficient || unjudged.valid() { + t.Fatal("unjudged outcome became an abstention") + } +} + +func TestInvalidClaimsNeverReachProvider(t *testing.T) { + for _, change := range []string{"statement", "batch key", "option key", "missing outcome", "unknown outcome"} { + t.Run(change, func(t *testing.T) { + claim := testClaim("Is this product running?") + key := "product" + switch change { + case "statement": + claim.Statement = " \n" + case "batch key": + key = " " + case "option key": + claim.Options[" "] = Option{Outcome: Holds} + case "missing outcome": + claim.Options["yes"] = Option{} + case "unknown outcome": + claim.Options["yes"] = Option{Outcome: "accepted"} + } + provider := answerProvider(func(map[string]Claim) map[string]Ruling { + t.Fatal("invalid claim reached provider") + return nil + }) + claims := map[string]Claim{key: claim} + if _, err := Judge(context.Background(), provider, "state", claims); err == nil { + t.Fatal("Judge accepted invalid claim") + } + if _, err := Cached(provider, 1).Judge(context.Background(), "state", claims); err == nil { + t.Fatal("cache accepted invalid claim") + } + }) + } +} + +const testInsufficient = "The evidence decides neither way." + +func testClaim(statement string) Claim { + return Claim{Statement: statement, Options: map[string]Option{"yes": {"", Holds}, "no": {"", Refuted}, OptionInsufficient: {testInsufficient, Insufficient}}} +} + +type answerProvider func(map[string]Claim) map[string]Ruling + +func (answerProvider) ID() string { return "answer-provider" } + +func (p answerProvider) Judge(_ context.Context, _ interface{}, questions map[string]Claim) (map[string]Ruling, error) { + return p(questions), nil +} + +func TestExplicitAbstentionAndCompleteContract(t *testing.T) { + claim := testClaim("?") + for _, confidence := range []float64{-0.1, 1.1, math.NaN(), math.Inf(1)} { + p := answerProvider(func(map[string]Claim) map[string]Ruling { + return map[string]Ruling{"c": {Option: "yes", Confidence: confidence}} + }) + if _, e := Judge(context.Background(), p, "state", map[string]Claim{"c": claim}); e == nil { + t.Fatalf("accepted confidence %v", confidence) + } + } + for _, conflict := range []bool{false, true} { + c := testClaim("?") + if conflict { + c.Options[OptionInsufficient] = Option{Outcome: Holds} + } else { + delete(c.Options, OptionInsufficient) + } + called := false + p := answerProvider(func(map[string]Claim) map[string]Ruling { called = true; return nil }) + if _, e := Judge(context.Background(), p, nil, map[string]Claim{"c": c}); e == nil || called { + t.Fatal("invalid abstention reached provider") + } + } + if r, e := Judge(context.Background(), nil, nil, nil); e != nil || len(r) != 0 { + t.Fatal("empty batch required a provider") + } + if _, e := Judge(context.Background(), nil, nil, map[string]Claim{"c": claim}); e == nil { + t.Fatal("missing provider accepted") + } +} diff --git a/jev/client.go b/jev/client.go new file mode 100644 index 0000000..6c4d790 --- /dev/null +++ b/jev/client.go @@ -0,0 +1,160 @@ +package jev + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "os" + "sync/atomic" + "time" +) + +const ( + DefaultEndpoint = "https://api.typesafe.ai/v1/systemone" + // DefaultMinConfidence is the threshold chosen for the pinned Jev model. + DefaultMinConfidence = 0.3 + DefaultModel = "jev-1.13.0" + EnvAPIKey = "TYPESAFE_API_KEY" +) + +// Client is the Provider that sends claims to the Jev HTTP API. +type Client struct { + Endpoint string + Model string + APIKey string + HTTP *http.Client + MaxRetries int // on 429 / 529, with exponential backoff + + // InputTokens sent so far; read with atomic.LoadInt64. + InputTokens int64 +} + +// NewClient reads the API key from TYPESAFE_API_KEY when apiKey is empty. +func NewClient(apiKey string) (*Client, error) { + if apiKey == "" { + apiKey = os.Getenv(EnvAPIKey) + } + if apiKey == "" { + return nil, fmt.Errorf("%s is not set", EnvAPIKey) + } + return &Client{ + Endpoint: DefaultEndpoint, + Model: DefaultModel, + APIKey: apiKey, + HTTP: &http.Client{Timeout: 30 * time.Second}, + MaxRetries: 3, + }, nil +} + +func (p *Client) ID() string { return "jev/" + p.Endpoint + "/" + p.Model } + +// wireClaim is only the Jev HTTP representation. Core Claim has no wire policy. +type wireClaim struct { + Type string `json:"type"` + Instructions string `json:"instructions"` + Criteria map[string]*string `json:"criteria"` +} + +func toWire(q Claim) wireClaim { + opts := make(map[string]*string, len(q.Options)) + for k, v := range q.Options { + if v.Description == "" { + opts[k] = nil // Jev's "no description" + } else { + description := v.Description + opts[k] = &description + } + } + return wireClaim{Type: "choice", Instructions: q.Statement, Criteria: opts} +} + +type APIError struct { + Status int + Body string +} + +func (e *APIError) Error() string { return fmt.Sprintf("typesafe: http %d: %s", e.Status, e.Body) } + +func (p *Client) Judge(ctx context.Context, state interface{}, claims map[string]Claim) (map[string]Ruling, error) { + if err := ValidateClaims(claims); err != nil { + return nil, err + } + req := struct { + State interface{} `json:"state"` + Model string `json:"model"` + Questions map[string]wireClaim `json:"questions"` + }{State: state, Model: p.Model, Questions: make(map[string]wireClaim, len(claims))} + for id, claim := range claims { + req.Questions[id] = toWire(claim) + } + body, err := json.Marshal(req) + if err != nil { + return nil, err + } + backoff := time.Second + for attempt := 0; ; attempt++ { + rulings, err := p.do(ctx, body) + if err == nil { + if err := ValidateRulings(claims, rulings); err != nil { + return nil, err + } + return rulings, nil + } + apiErr, ok := err.(*APIError) + if !ok || (apiErr.Status != 429 && apiErr.Status != 529) || attempt >= p.MaxRetries { + return nil, err + } + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-time.After(backoff): + } + backoff *= 2 + } +} + +func (p *Client) do(ctx context.Context, body []byte) (map[string]Ruling, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodPost, p.Endpoint, bytes.NewReader(body)) + if err != nil { + return nil, err + } + req.Header.Set("Authorization", "Bearer "+p.APIKey) + req.Header.Set("Content-Type", "application/json") + + resp, err := p.HTTP.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + raw, err := io.ReadAll(resp.Body) + if err != nil { + return nil, err + } + if resp.StatusCode != http.StatusOK { + return nil, &APIError{Status: resp.StatusCode, Body: string(raw)} + } + var out struct { + Answers map[string]struct { + Choice string `json:"choice"` + Confidence *float64 `json:"confidence"` + } `json:"answers"` + Usage struct { + InputTokens int64 `json:"input_tokens"` + } `json:"usage"` + } + if err := json.Unmarshal(raw, &out); err != nil { + return nil, fmt.Errorf("typesafe: decode response: %w", err) + } + atomic.AddInt64(&p.InputTokens, out.Usage.InputTokens) + rulings := make(map[string]Ruling, len(out.Answers)) + for id, answer := range out.Answers { + if answer.Confidence == nil { + return nil, fmt.Errorf("typesafe: claim %q omitted confidence", id) + } + rulings[id] = Ruling{Option: answer.Choice, Confidence: *answer.Confidence} + } + return rulings, nil +} diff --git a/jev/client_test.go b/jev/client_test.go new file mode 100644 index 0000000..65565db --- /dev/null +++ b/jev/client_test.go @@ -0,0 +1,90 @@ +package jev + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" +) + +func TestClientWireFormatAndRetry(t *testing.T) { + calls := 0 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls++ + if r.Header.Get("Authorization") != "Bearer k" { + t.Errorf("auth header: %q", r.Header.Get("Authorization")) + } + if calls == 1 { + w.WriteHeader(529) + return + } + var req struct { + Model string + Questions map[string]map[string]interface{} + } + json.NewDecoder(r.Body).Decode(&req) + c := req.Questions["c"] + criteria, _ := c["criteria"].(map[string]interface{}) + if req.Model != DefaultModel || c["type"] != "choice" || criteria["x"] != nil || criteria["y"] != "why" || criteria["z"] != "another description" { + t.Errorf("request: %+v", req) + } + w.Write([]byte(`{"model":"jev-1.13.0","answers":{"c":{"type":"choice","choice":"y","confidence":0.8}},"usage":{"input_tokens":10,"output_tokens":1}}`)) + })) + defer srv.Close() + + p, _ := NewClient("k") + p.Endpoint = srv.URL + answers, err := p.Judge(context.Background(), "state", map[string]Claim{ + "c": {Statement: "?", Options: map[string]Option{"x": {Outcome: Refuted}, "y": {Description: "why", Outcome: Holds}, "z": {Description: "another description", Outcome: Holds}, OptionInsufficient: {Outcome: Insufficient}}}, + }) + if err != nil { + t.Fatal(err) + } + if calls != 2 || answers["c"].Option != "y" || p.InputTokens != 10 { + t.Fatalf("calls=%d answers=%+v tokens=%d", calls, answers, p.InputTokens) + } + if p.ID() != "jev/"+srv.URL+"/jev-1.13.0" { + t.Fatalf("id %s", p.ID()) + } +} + +func TestClientRequiresExplicitWireConfidence(t *testing.T) { + for _, tc := range []struct { + name, answer string + valid bool + }{ + {"missing", `{"choice":"yes"}`, false}, + {"null", `{"choice":"yes","confidence":null}`, false}, + {"zero", `{"choice":"yes","confidence":0}`, true}, + {"unknown option", `{"choice":"other","confidence":0.9}`, false}, + {"core format", `{"option":"yes","confidence":0.9}`, false}, + } { + t.Run(tc.name, func(t *testing.T) { + calls := 0 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls++ + w.Write([]byte(`{"answers":{"c":` + tc.answer + `}}`)) + })) + defer srv.Close() + p, _ := NewClient("local-test") + p.Endpoint = srv.URL + claims := map[string]Claim{"c": {Statement: "Is it running?", Options: map[string]Option{ + "yes": {Outcome: Holds}, OptionInsufficient: {Outcome: Insufficient}, + }}} + rulings, err := p.Judge(context.Background(), "state", claims) + if (err == nil) != tc.valid || calls != 1 { + t.Fatalf("rulings=%v error=%v calls=%d", rulings, err, calls) + } + if tc.valid && (rulings["c"].Option != "yes" || rulings["c"].Confidence != 0) { + t.Fatalf("zero confidence changed: %v", rulings) + } + invalid := claims["c"] + invalid.Statement = "" + claims["c"] = invalid + if _, err := p.Judge(context.Background(), "state", claims); err == nil || calls != 1 { + t.Fatal("invalid claim reached HTTP boundary") + } + }) + } +} diff --git a/jev/doc.go b/jev/doc.go new file mode 100644 index 0000000..fcb1f9a --- /dev/null +++ b/jev/doc.go @@ -0,0 +1,16 @@ +// Package jev asks a model to rule on claims about shared evidence. +// +// A Claim is a statement with a finite set of options; each option declares +// what choosing it means (Holds, Refuted or Insufficient), and every claim +// offers an explicit "insufficient" abstention. A Provider picks one option +// per claim with a confidence; Claim.Resolve is the only interpretation of +// that pick. The model never writes free text back. +// +// c, err := jev.NewClient("") // TYPESAFE_API_KEY +// p := jev.Cached(c, jev.DefaultCacheSize) +// rulings, err := jev.Judge(ctx, p, state, claims) +// outcome := claims["k"].Resolve(rulings["k"], jev.DefaultMinConfidence) +// +// Client speaks the TypeSafe Jev (System One) HTTP API +// (https://docs.typesafe.ai/api.md). Other models plug in as a Provider. +package jev diff --git a/jev/go.mod b/jev/go.mod new file mode 100644 index 0000000..5bbb805 --- /dev/null +++ b/jev/go.mod @@ -0,0 +1,3 @@ +module github.com/chainreactors/utils/jev + +go 1.17 diff --git a/jev/provider.go b/jev/provider.go new file mode 100644 index 0000000..4719930 --- /dev/null +++ b/jev/provider.go @@ -0,0 +1,54 @@ +package jev + +import ( + "context" + "encoding/json" + "fmt" +) + +// Provider evaluates finite-option claims against shared evidence. It returns +// the original selections; Claim.Resolve alone interprets their meaning. +// Implementations must not mutate evidence, claims or their options. +type Provider interface { + // ID identifies the endpoint and model. Change it when answers may change. + ID() string + Judge(context.Context, interface{}, map[string]Claim) (map[string]Ruling, error) +} + +// Judge validates a batch, freezes the evidence as JSON once so every provider +// sees the same bytes, and validates the complete answer. A failed batch +// returns no rulings. Resolve each ruling against its claim before acting. +func Judge(ctx context.Context, p Provider, state interface{}, claims map[string]Claim) (map[string]Ruling, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + if err := ValidateClaims(claims); err != nil { + return nil, err + } + if len(claims) == 0 { + return map[string]Ruling{}, nil + } + if p == nil { + return nil, fmt.Errorf("jev: no provider") + } + data, err := json.Marshal(state) + if err != nil { + return nil, err + } + got, err := p.Judge(ctx, json.RawMessage(data), claims) + if err != nil { + return nil, err + } + if err := ValidateRulings(claims, got); err != nil { + return nil, err + } + return copyRulings(got), nil +} + +func copyRulings(in map[string]Ruling) map[string]Ruling { + out := make(map[string]Ruling, len(in)) + for key, r := range in { + out[key] = r + } + return out +}