Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 45 additions & 0 deletions jev/README.md
Original file line number Diff line number Diff line change
@@ -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 会报错。
101 changes: 101 additions & 0 deletions jev/cache.go
Original file line number Diff line number Diff line change
@@ -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
}
179 changes: 179 additions & 0 deletions jev/cache_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
Loading