diff --git a/cmd/engine/main.go b/cmd/engine/main.go index 352267c..9eed587 100644 --- a/cmd/engine/main.go +++ b/cmd/engine/main.go @@ -19,6 +19,8 @@ import ( "time" "github.com/hallelx2/llmgate" + "github.com/hallelx2/llmgate/judge/typesafe" + "github.com/hallelx2/llmgate/middleware/retry" "github.com/hallelx2/llmgate/pricing" "github.com/hallelx2/llmgate/provider/anthropic" "github.com/hallelx2/llmgate/provider/gemini" @@ -203,10 +205,22 @@ func run() error { ) } + judge, err := buildJudge(cfg.LLM.Judge) + if err != nil { + logger.Error("judge: config invalid", "err", err) + os.Exit(1) + } + if judge != nil { + logger.Info("judge: typesafe enabled — contents-page detection and page resolution run on the Judge") + } else { + logger.Warn("judge: none configured — TOC judgements run on the generative driver, one call per page (set TYPESAFE_API_KEY)") + } pipeline := ingest.NewPipeline(ingest.Pipeline{ DB: pool, Storage: store, LLM: llmClient, + Judge: judge, + JudgeThreshold: cfg.LLM.Judge.Threshold, Parsers: ingest.RegistryFromIngestParams(tableOptsFromConfig(cfg.Ingest.Tables), cfg.Ingest.MaxSections, time.Duration(cfg.Ingest.ParseTimeoutSeconds)*time.Second), Logger: logger, Mode: cfg.Ingest.Mode, @@ -396,6 +410,26 @@ func modelFor(c config.LLMConfig) string { return "" } +// buildJudge returns the configured Judge, or nil when llm.judge has +// no API key. Retries share the same schedule as every other provider +// call; a Judge request that fails past them is handled by the TOC +// builder, which keeps extraction's pages rather than degrading +// silently (HAL-1369). +func buildJudge(c config.JudgeBlock) (llmgate.Judge, error) { + if c.TypeSafe.APIKey == "" { + return nil, nil + } + j, err := typesafe.New(typesafe.Config{ + APIKey: c.TypeSafe.APIKey, + BaseURL: c.TypeSafe.BaseURL, + Model: c.TypeSafe.Model, + }) + if err != nil { + return nil, err + } + return retry.NewJudge(retry.Config{MaxRetries: 3})(j), nil +} + func buildLLM(c config.LLMConfig) (llmgate.Client, error) { switch c.Driver { case "anthropic": diff --git a/cmd/server/main.go b/cmd/server/main.go index 129fbeb..13c1089 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -27,6 +27,8 @@ import ( "time" "github.com/hallelx2/llmgate" + "github.com/hallelx2/llmgate/judge/typesafe" + "github.com/hallelx2/llmgate/middleware/retry" "github.com/hallelx2/llmgate/pricing" "github.com/hallelx2/llmgate/provider/anthropic" "github.com/hallelx2/llmgate/provider/gemini" @@ -201,10 +203,22 @@ func run() error { } // ── Ingest pipeline ─────────────────────────────────────────── + judge, err := buildJudge(cfg.Engine.LLM.Judge) + if err != nil { + logger.Error("judge: config invalid", "err", err) + os.Exit(1) + } + if judge != nil { + logger.Info("judge: typesafe enabled — contents-page detection and page resolution run on the Judge") + } else { + logger.Warn("judge: none configured — TOC judgements run on the generative driver, one call per page (set TYPESAFE_API_KEY)") + } pipeline := ingest.NewPipeline(ingest.Pipeline{ DB: pool, Storage: store, LLM: llmClient, + Judge: judge, + JudgeThreshold: cfg.Engine.LLM.Judge.Threshold, Parsers: ingest.RegistryFromIngestParams(tableOptsFromConfig(cfg.Engine.Ingest.Tables), cfg.Engine.Ingest.MaxSections, time.Duration(cfg.Engine.Ingest.ParseTimeoutSeconds)*time.Second), Logger: logger, Mode: cfg.Engine.Ingest.Mode, @@ -396,6 +410,26 @@ func modelFor(c enginecfg.LLMConfig) string { return "" } +// buildJudge returns the configured Judge, or nil when llm.judge has +// no API key. Retries share the same schedule as every other provider +// call; a Judge request that fails past them is handled by the TOC +// builder, which keeps extraction's pages rather than degrading +// silently (HAL-1369). +func buildJudge(c enginecfg.JudgeBlock) (llmgate.Judge, error) { + if c.TypeSafe.APIKey == "" { + return nil, nil + } + j, err := typesafe.New(typesafe.Config{ + APIKey: c.TypeSafe.APIKey, + BaseURL: c.TypeSafe.BaseURL, + Model: c.TypeSafe.Model, + }) + if err != nil { + return nil, err + } + return retry.NewJudge(retry.Config{MaxRetries: 3})(j), nil +} + func buildLLM(c enginecfg.LLMConfig) (llmgate.Client, error) { switch c.Driver { case "anthropic": diff --git a/config.example.yaml b/config.example.yaml index 61a492b..c19c523 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -107,6 +107,20 @@ llm: model: "gemini-2.0-flash" reasoning_model: "gemini-2.5-pro" + # Judge: a System One model (TypeSafe Jev) that answers the ingest + # pipeline's judgements — contents-page detection and page resolution — + # in one batched request per document instead of a generative call per + # page. Enabled exactly when an api_key is present; also read from + # VLE_TYPESAFE_API_KEY or TYPESAFE_API_KEY. Without it, long filings + # lose every leaf's page (HAL-1367) and the TOC stage runs minutes + # slower on the driver above. + judge: + typesafe: + api_key: "" + # base_url: "" # override the endpoint; empty = api.typesafe.ai + # model: "" # override the model alias + # threshold: 0.5 # Noul probability that counts as "yes" + retrieval: # strategy: single-pass | chunked-tree | agentic | treewalk # diff --git a/config.server.example.yaml b/config.server.example.yaml index 02d7835..260935f 100644 --- a/config.server.example.yaml +++ b/config.server.example.yaml @@ -90,6 +90,13 @@ engine: # api_key: "" # model: "gemini-2.0-flash" # reasoning_model: "" + # Judge (TypeSafe Jev): batched contents-page detection and page + # resolution during ingest. Enabled when api_key is set — also read + # from VLE_TYPESAFE_API_KEY / TYPESAFE_API_KEY. See HAL-1367. + judge: + typesafe: + api_key: "" + # threshold: 0.5 retrieval: strategy: "chunked-tree" # "single-pass" or "chunked-tree" diff --git a/pkg/config/config.go b/pkg/config/config.go index beaefee..238c878 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -376,6 +376,32 @@ type LLMConfig struct { Anthropic AnthropicBlock `yaml:"anthropic"` OpenAI OpenAIBlock `yaml:"openai"` Gemini GeminiBlock `yaml:"gemini"` + + // Judge configures the System One model that answers the pipeline's + // judgements — contents-page detection and page resolution — in one + // batched request each, instead of a generative call per page. Left + // empty, those steps run on the generative driver above, one call at + // a time, and a 10-K's leaves lose their pages (HAL-1367). + Judge JudgeBlock `yaml:"judge"` +} + +// JudgeBlock configures the Judge. Only TypeSafe is supported today; the +// Judge is enabled exactly when an API key is present. +type JudgeBlock struct { + TypeSafe TypeSafeBlock `yaml:"typesafe"` + + // Threshold is the Noul probability at or above which a Judge answer + // counts as yes. Zero selects the builder's default (0.5). + Threshold float64 `yaml:"threshold"` +} + +// TypeSafeBlock configures the TypeSafe System One provider. +type TypeSafeBlock struct { + APIKey string `yaml:"api_key"` + // BaseURL overrides the endpoint. Empty = api.typesafe.ai. + BaseURL string `yaml:"base_url"` + // Model overrides the model alias. Empty = the provider's default. + Model string `yaml:"model"` } // AnthropicBlock configures the Anthropic provider. @@ -903,6 +929,13 @@ func applyEnvOverrides(c *Config) { if v := os.Getenv("VLE_GEMINI_API_KEY"); v != "" { c.LLM.Gemini.APIKey = v } + // The Judge key accepts the deploy layer's VLS_ prefix and the bare + // TYPESAFE_API_KEY — what the llmgate live tests and the bench + // commands already read — so one export enables it everywhere. + // VLE_-prefixed wins if several are set. + if v := firstEnv("VLE_TYPESAFE_API_KEY", "VLS_TYPESAFE_API_KEY", "TYPESAFE_API_KEY"); v != "" { + c.LLM.Judge.TypeSafe.APIKey = v + } // Accept both VLE_-prefixed and bare QSTASH_* env vars. The bare // names match what the Upstash console documents and what the // dashboard already uses, so ops can set them once for both diff --git a/pkg/config/judge_config_test.go b/pkg/config/judge_config_test.go new file mode 100644 index 0000000..e708e0e --- /dev/null +++ b/pkg/config/judge_config_test.go @@ -0,0 +1,29 @@ +package config + +import "testing" + +func TestJudgeKeyFromEnv(t *testing.T) { + t.Setenv("VLE_TYPESAFE_API_KEY", "") + t.Setenv("TYPESAFE_API_KEY", "bare") + c := Default() + applyEnvOverrides(&c) + if c.LLM.Judge.TypeSafe.APIKey != "bare" { + t.Errorf("bare TYPESAFE_API_KEY not picked up: %q", c.LLM.Judge.TypeSafe.APIKey) + } + t.Setenv("VLE_TYPESAFE_API_KEY", "prefixed") + c = Default() + applyEnvOverrides(&c) + if c.LLM.Judge.TypeSafe.APIKey != "prefixed" { + t.Errorf("VLE_ prefix should win: %q", c.LLM.Judge.TypeSafe.APIKey) + } +} + +func TestJudgeIsOffByDefault(t *testing.T) { + t.Setenv("VLE_TYPESAFE_API_KEY", "") + t.Setenv("TYPESAFE_API_KEY", "") + c := Default() + applyEnvOverrides(&c) + if c.LLM.Judge.TypeSafe.APIKey != "" { + t.Errorf("no key configured, got %q", c.LLM.Judge.TypeSafe.APIKey) + } +} diff --git a/pkg/ingest/ingest.go b/pkg/ingest/ingest.go index 3427099..9559fb5 100644 --- a/pkg/ingest/ingest.go +++ b/pkg/ingest/ingest.go @@ -170,6 +170,18 @@ type Pipeline struct { // per-stage semaphore). Default applied by NewPipeline: 12. GlobalLLMConcurrency int + // Judge, when non-nil, answers the TOC stage's judgements — + // contents-page detection and page resolution — in one batched + // request each. Without it those steps run on LLM, one generative + // call per page, and long filings lose every leaf's page (HAL-1367). + // Wired from config llm.judge by cmd/server and cmd/engine; nil in + // Pipeline literals that do not set it, which keeps the old path. + Judge llmgate.Judge + + // JudgeThreshold is the Noul probability at or above which a Judge + // answer counts as yes. Zero selects TOCBuilder's default. + JudgeThreshold float64 + // TOCEnabled toggles the LLM-built table-of-contents stage. The // stage runs after summarize+HyDE on PDF inputs and persists the // resulting tree on documents.toc_tree (JSONB). Failures are @@ -460,6 +472,12 @@ func (p *Pipeline) runTOCBuilder(ctx context.Context, docID tree.DocumentID, par Concurrency: p.TOCConcurrency, TOCCheckPages: p.TOCCheckPages, LLMCallTimeout: p.LLMCallTimeout, + Judge: p.Judge, + JudgeThreshold: p.JudgeThreshold, + // The minimum-context path: prefilter, truncation, two-stage + // scan. Measured on FinanceBench as the fastest detection that + // lost nothing (HAL-1366). + MinimalContext: p.Judge != nil, } nodes, usage, err := builder.Build(ctx, pages) if err != nil { @@ -470,7 +488,13 @@ func (p *Pipeline) runTOCBuilder(ctx context.Context, docID tree.DocumentID, par "llm_calls", usage.LLMCalls, "input_tokens", usage.InputTokens, "output_tokens", usage.OutputTokens, + "judge", p.Judge != nil, ) + // A degraded build is not a failed one, but it is not what was + // configured either, and it must not look like success in the log. + for _, d := range usage.Degraded { + log.Warn("ingest: toc-builder degraded", "step", d) + } if len(nodes) == 0 { return nil } diff --git a/pkg/ingest/toc_builder.go b/pkg/ingest/toc_builder.go index 1c8fe06..574ebd4 100644 --- a/pkg/ingest/toc_builder.go +++ b/pkg/ingest/toc_builder.go @@ -136,6 +136,19 @@ type Usage struct { TotalTokens int CostUSD float64 LLMCalls int + + // Degraded lists the Judge-path steps that could not complete and + // what Build did instead. Empty means every step ran as configured. + // A document built with a non-empty Degraded is not wrong, but it is + // not what was asked for, and the caller must be able to see that: + // VERIZON_2022_10K once ingested with no page on any leaf after a + // single failed Judge request, and reported success (HAL-1369). + Degraded []string +} + +// degrade records a Judge-path step that fell back. +func (u *Usage) degrade(step, what string) { + u.Degraded = append(u.Degraded, step+": "+what) } // add folds the per-response usage from one LLM call into the @@ -210,10 +223,16 @@ func (b *TOCBuilder) Build(ctx context.Context, pages []PageText) ([]tree.TOCNod // verification with search, and it is what makes the extraction body // window irrelevant to page accuracy (HAL-1367). Plain verification // remains the fallback when resolution cannot run. - if resolved, handled := b.resolvePagesJudge(ctx, nodes, pages, tocPages, &usage); handled { - applyResolvedPages(nodes, resolved) - } else if verdicts, handled := b.verifyTitlesJudge(ctx, nodes, pages, &usage); handled { - applyJudgeVerdicts(nodes, verdicts) + // + // A failed Judge request is not "no Judge". The generative verifier + // asks about extraction's printed page numbers, which are wrong for + // every leaf past the cover of a long filing, so falling to it after + // a transport failure zeroes the tree and looks like success + // (HAL-1369). With a Judge configured, resolution is retried once + // with a fresh budget; if it still fails, extraction's pages are + // kept as they are and the degradation is recorded on Usage. + if b.Judge != nil { + b.resolvePagesOrKeep(ctx, nodes, pages, tocPages, &usage) } else { b.verifyTitlesConcurrent(ctx, nodes, pages, concurrency, &usage) } @@ -229,6 +248,36 @@ func (b *TOCBuilder) Build(ctx context.Context, pages []PageText) ([]tree.TOCNod return nodes, usage, nil } +// resolverAttempts is how many times Build asks the Judge to resolve +// pages before keeping extraction's. Two: the first failure is almost +// always transport, and a resolver batch is two cheap requests. +const resolverAttempts = 2 + +// resolvePagesOrKeep runs Judge page resolution with one retry. On +// exhaustion it leaves the tree exactly as extraction produced it and +// records the fact; it never routes a Judge-path document through the +// generative verifier. +func (b *TOCBuilder) resolvePagesOrKeep(ctx context.Context, nodes []tree.TOCNode, pages []PageText, exclude []int, usage *Usage) { + var lastErr error + for attempt := 1; attempt <= resolverAttempts; attempt++ { + resolved, handled, err := b.resolvePagesJudgeErr(ctx, nodes, pages, exclude, usage) + if err == nil { + if handled { + applyResolvedPages(nodes, resolved) + } + // handled=false with no error means there was nothing to + // resolve (no leaves with titles); extraction's pages stand. + return + } + lastErr = err + log.Printf("toc: judge page resolution attempt %d/%d failed: %v", attempt, resolverAttempts, err) + if ctx.Err() != nil { + break + } + } + usage.degrade("page resolution", fmt.Sprintf("kept extraction's pages after %d failed Judge attempts: %v", resolverAttempts, lastErr)) +} + // detectTOCPages scans the first tocCheck pages with the // TreeWalk-style single-page detector. Returns the 1-indexed page // numbers (in order) the LLM judged as table-of-contents pages. diff --git a/pkg/ingest/toc_builder_test.go b/pkg/ingest/toc_builder_test.go index 81e5ba9..68733b8 100644 --- a/pkg/ingest/toc_builder_test.go +++ b/pkg/ingest/toc_builder_test.go @@ -2,6 +2,7 @@ package ingest import ( "context" + "errors" "strings" "sync" "sync/atomic" @@ -527,3 +528,115 @@ func TestSynthetic10KFourTopLevelNodes(t *testing.T) { } } } + +// judgeFailingFirst is a Judge whose page-resolution batches fail the +// first n times and answer yes afterwards; detection always says no so +// the scripted LLM's no-TOC extractor supplies the tree. +func judgeFailingFirst(n int) (*llmgate.MockJudge, *atomic.Int32) { + var resolverCalls atomic.Int32 + j := &llmgate.MockJudge{ + Respond: func(_ context.Context, req llmgate.JudgeRequest) (*llmgate.Judgment, error) { + ans := make(map[string]llmgate.Answer, len(req.Questions)) + isResolver := false + for id := range req.Questions { + if strings.HasPrefix(id, "r_") { + isResolver = true + } + } + if isResolver { + if int(resolverCalls.Add(1)) <= n { + return nil, errors.New("typesafe: request failed") + } + st, _ := req.State.(map[string]any) + for id := range req.Questions { + p := 0.05 + if item, ok := st[id].(map[string]any); ok { + title, _ := item["title"].(string) + excerpt, _ := item["excerpt"].(string) + if strings.HasPrefix(excerpt, title) { + p = 0.95 + } + } + ans[id] = llmgate.NoulAnswer{Noul: p} + } + } else { + for id := range req.Questions { + ans[id] = llmgate.NoulAnswer{Noul: 0.0} + } + } + return &llmgate.Judgment{Model: "mock", Answers: ans, + Usage: llmgate.Usage{InputTokens: 10, TotalTokens: 10, TokensReported: true}}, nil + }, + } + return j, &resolverCalls +} + +func buildWithJudge(t *testing.T, j llmgate.Judge) ([]tree.TOCNode, Usage, *atomic.Int32) { + t.Helper() + llm := &scriptedLLM{} + var verifierCalls atomic.Int32 + llm.route = func(prompt string) string { + switch { + case strings.Contains(prompt, "hierarchical tree structure"): + // Extraction guesses the printed page numbers, which are wrong. + return `{"nodes":[ + {"structure":"1","title":"Item 1. Business","physical_index":""}, + {"structure":"2","title":"Item 1A. Risk Factors","physical_index":""} + ]}` + case strings.Contains(prompt, "section starts at the beginning"): + verifierCalls.Add(1) + return `{"start_begin":"no"}` + } + return `{"toc_detected":"no"}` + } + pages := []PageText{ + {PageNumber: 1, Text: "Cover page."}, + {PageNumber: 3, Text: "Item 1. Business\nWe make things."}, + {PageNumber: 8, Text: "Item 1A. Risk Factors\nThings go wrong."}, + } + b := &TOCBuilder{LLM: llm, Judge: j, TOCCheckPages: 3, Concurrency: 2} + nodes, usage, err := b.Build(context.Background(), pages) + if err != nil { + t.Fatalf("Build: %v", err) + } + return nodes, usage, &verifierCalls +} + +// One failed Judge request is transport, not absence: Build retries the +// resolver and the document gets its real pages. +func TestBuildRetriesTheResolverOnce(t *testing.T) { + j, calls := judgeFailingFirst(1) + nodes, usage, verifier := buildWithJudge(t, j) + if calls.Load() != 2 { + t.Errorf("resolver batches: got %d want 2 (one failure, one success)", calls.Load()) + } + if nodes[0].StartPage != 3 || nodes[1].StartPage != 8 { + t.Errorf("pages not resolved after the retry: %+v", nodes) + } + if len(usage.Degraded) != 0 { + t.Errorf("a recovered retry is not a degradation: %v", usage.Degraded) + } + if verifier.Load() != 0 { + t.Errorf("the generative verifier ran %d times on a Judge-path document", verifier.Load()) + } +} + +// After the retries are spent, extraction's pages are kept, the +// generative verifier — which would reject every one of them — is never +// consulted, and the caller can see the degradation. +func TestBuildKeepsClaimedPagesWhenTheResolverIsDown(t *testing.T) { + j, calls := judgeFailingFirst(1000) + nodes, usage, verifier := buildWithJudge(t, j) + if calls.Load() != resolverAttempts { + t.Errorf("resolver batches: got %d want %d", calls.Load(), resolverAttempts) + } + if nodes[0].StartPage != 1 || nodes[1].StartPage != 6 { + t.Errorf("extraction's pages should be kept untouched, got %+v", nodes) + } + if verifier.Load() != 0 { + t.Errorf("the generative verifier ran %d times; it would have zeroed the tree", verifier.Load()) + } + if len(usage.Degraded) != 1 || !strings.Contains(usage.Degraded[0], "page resolution") { + t.Errorf("degradation not recorded: %v", usage.Degraded) + } +}