Skip to content
Open
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
5 changes: 3 additions & 2 deletions cmd/navbench/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -73,8 +73,9 @@ func main() {
pdfs := flag.String("pdfs", "", "directory of filings")
out := flag.String("out", "", "JSONL of per-question outcomes")
maxLeaves := flag.Int("leaves", 0, "cap on sections read per question; 0 lets the page budget decide")
maxPages := flag.Int("pages", 40, "pages judged per question")
maxPages := flag.Int("pages", 0, "pages judged in full per question (0 = the engine default)")
limitQ := flag.Int("limit", 0, "stop after this many questions (0 = all)")
skimAll := flag.Bool("skim-all", false, "skim every page's head alongside the section ranking (HAL-1566), as the persisted-pages path does")
parallel := flag.Int("parallel", 1, "questions in flight at once; the provider's adaptive limiter governs requests")
flag.Parse()
stats := &judgestats.Recorder{}
Expand All @@ -97,7 +98,7 @@ func main() {
lim := limit.New(limit.Config{Initial: 16, OnChange: func(e limit.Event) {
fmt.Fprintf(os.Stderr, " limiter %s %d -> %d %v\n", e.Cause, e.From, e.To, e.Err)
}})
nav := &retrieval.JudgeNavigator{Judge: retry.NewJudge(retry.Config{MaxRetries: 3})(limit.Judge(lim)(tj)), MaxLeaves: *maxLeaves, MaxPages: *maxPages}
nav := &retrieval.JudgeNavigator{Judge: retry.NewJudge(retry.Config{MaxRetries: 3})(limit.Judge(lim)(tj)), MaxLeaves: *maxLeaves, MaxPages: *maxPages, SkimAll: *skimAll}

qs := readQuestions(*qPath)
if *limitQ > 0 && len(qs) > *limitQ {
Expand Down
216 changes: 191 additions & 25 deletions pkg/retrieval/judgewalk.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,8 @@ type JudgeNavigator struct {
// five one-page notes and the budget went unused (HAL-1374).
MaxLeaves int

// MaxPages bounds how many pages are judged in full. Zero selects 40.
// MaxPages bounds how many pages are judged in full. Zero selects
// defaultNavMaxPages.
MaxPages int

// CoarsePages bounds how many gathered pages the coarse pass may
Expand Down Expand Up @@ -118,11 +119,46 @@ type JudgeNavigator struct {
// not more waiting — and the provider is far happier with many small
// requests than a few large ones (see defaultNavReqTokens).
RequestBudgetTokens int

// SkimAll skims the head of every page in the document at the same
// time as the section ranking, instead of waiting for the ranking and
// skimming only the chosen sections' pages (HAL-1566).
//
// The skim only ever depended on the ranking because of WHICH pages
// it read; the head score of a page does not depend on the ranking
// at all. Scoring every page's head up front makes the two requests
// independent, and one of navigation's four sequential levels goes
// away. Which pages are read in full is unchanged: the candidates
// are still the ranked sections' pages up to CoarsePages, ordered by
// head score.
//
// It costs the heads of the pages outside the chosen sections. Set it
// only where loading every leaf's pages is cheap — the persisted-
// pages path, where they are already in memory — never where each
// load is a storage read.
SkimAll bool

// SkimAllMaxPages is the largest document SkimAll applies to; a
// longer one falls back to the sequential skim, so the extra heads
// stay bounded. Zero selects defaultNavSkimAllMax.
SkimAllMaxPages int
}

const (
defaultNavThreshold = 0.5
defaultNavMaxPages = 40
// defaultNavMaxPages is how many pages are read in full, and full
// pages are ~60% of a query's tokens. Measured 2026-09-28 on 40
// FinanceBench questions with SkimAll, repeated runs:
//
// 40 pages hit 35/36/35 $0.00431/q
// 30 pages hit 36/36 $0.00388/q
// 20 pages hit 35/36 $0.00345/q
//
// Every row sits inside Jev's run-to-run noise (35–36). 30 keeps
// headroom for multi-page answers, which FinanceBench barely tests
// and a smaller read would hurt first; 20 is the frugal setting for
// single-fact corpora.
defaultNavMaxPages = 30
defaultNavCoarse = 120
defaultNavHeadChars = 700
defaultNavPageChars = 6000
Expand All @@ -143,6 +179,11 @@ const (
navLeafBatch = 120
navMinEvidencePages = 2
navLeafStateMaxChars = 300
// defaultNavSkimAllMax bounds SkimAll: a head is ~150 tokens, so 400
// pages is ~60k tokens a query, about $0.0025 at Jev's input price —
// against the ~18k the sequential skim of 120 heads already spends.
// FinanceBench's filings run 60–300 pages.
defaultNavSkimAllMax = 400
)

func (n *JudgeNavigator) threshold() float64 {
Expand Down Expand Up @@ -187,6 +228,13 @@ func (n *JudgeNavigator) pageChars() int {
return defaultNavPageChars
}

func (n *JudgeNavigator) skimAllMax() int {
if n.SkimAllMaxPages > 0 {
return n.SkimAllMaxPages
}
return defaultNavSkimAllMax
}

func (n *JudgeNavigator) reqTokens() int {
if n.RequestBudgetTokens > 0 {
return n.RequestBudgetTokens
Expand All @@ -206,7 +254,11 @@ func (n *JudgeNavigator) RankLeaves(ctx context.Context, query string, leaves []
for i, l := range leaves {
scores[i] = LeafScore{Leaf: l}
}
requests := 0
type batch struct {
state map[string]any
questions map[string]llmgate.Question
}
var batches []batch
for start := 0; start < len(leaves); start += navLeafBatch {
end := start + navLeafBatch
if end > len(leaves) {
Expand Down Expand Up @@ -243,21 +295,49 @@ func (n *JudgeNavigator) RankLeaves(ctx context.Context, query string, leaves []
},
}
}
res, err := n.Judge.Judge(ctx, llmgate.JudgeRequest{State: state, Questions: questions})
if err != nil {
return nil, usage, requests, err
}
requests++
usage.Add(judgeUsage(res))
for qk := range questions {
p, err := res.Noul(qk)
batches = append(batches, batch{state, questions})
}

// A tree with more than one batch of leaves sends them together, as
// rankPages does; they share nothing but the question.
ctx, cancel := context.WithCancel(ctx)
defer cancel()
var (
wg sync.WaitGroup
mu sync.Mutex
firstErr error
requests int
)
for _, b := range batches {
wg.Add(1)
go func(b batch) {
defer wg.Done()
res, err := n.Judge.Judge(ctx, llmgate.JudgeRequest{State: b.state, Questions: b.questions})
mu.Lock()
defer mu.Unlock()
if err != nil {
continue
if firstErr == nil {
firstErr = err
cancel()
}
return
}
var i int
fmt.Sscanf(qk, "l_%d", &i)
scores[i].P = p
}
requests++
usage.Add(judgeUsage(res))
for qk := range b.questions {
p, err := res.Noul(qk)
if err != nil {
continue
}
var i int
fmt.Sscanf(qk, "l_%d", &i)
scores[i].P = p
}
}(b)
}
wg.Wait()
if firstErr != nil {
return nil, usage, requests, firstErr
}
sort.SliceStable(scores, func(i, j int) bool { return scores[i].P > scores[j].P })
return scores, usage, requests, nil
Expand Down Expand Up @@ -389,12 +469,63 @@ func (n *JudgeNavigator) Navigate(ctx context.Context, query string, leaves []Na
if len(leaves) == 0 {
return out, nil
}
ranked, u, r, err := n.RankLeaves(ctx, query, leaves)
if err != nil {
return nil, fmt.Errorf("judgewalk: rank leaves: %w", err)
maxP := n.maxPages()

// With SkimAll the section ranking and the head skim of every page go
// out together; otherwise the skim waits for the ranking below.
var (
ranked []LeafScore
heads []PageScore
)
skimmed := false
if n.SkimAll {
all, err := allLeafPages(ctx, leaves, loadPages)
if err != nil {
return nil, err
}
if len(all) > maxP && len(all) <= n.skimAllMax() {
skimmed = true
var (
wg sync.WaitGroup
lu, hu Usage
lr, hr int
rankErr, headErr error
)
wg.Add(2)
go func() {
defer wg.Done()
ranked, lu, lr, rankErr = n.RankLeaves(ctx, query, leaves)
}()
go func() {
defer wg.Done()
heads, hu, hr, headErr = n.rankPages(ctx, query, all, n.headChars(), true)
}()
wg.Wait()
Comment on lines +495 to +503

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

issue (bug_risk): When either the concurrent section-ranking or head-skimming operation fails, the other operation keeps running because each ranking method creates and owns a separate child context; Navigate only waits for both and never cancels the sibling. The remaining requests therefore continue consuming provider capacity and tokens after the navigation has already become unrecoverable.

Triggers: When one of the two concurrent ranking operations returns an error while the other still has in-flight requests.

Suggested fix: Create one cancellable child context in Navigate and pass it to both operations, or use an errgroup.Group so either error cancels the sibling immediately.

if rankErr != nil {
return nil, fmt.Errorf("judgewalk: rank leaves: %w", rankErr)
}
if headErr != nil {
return nil, fmt.Errorf("judgewalk: rank page heads: %w", headErr)
}
out.Usage.Add(lu)
out.Usage.Add(hu)
out.Requests += lr + hr
out.Coarse = heads
}
}
if ranked == nil {
var (
u Usage
r int
err error
)
ranked, u, r, err = n.RankLeaves(ctx, query, leaves)
if err != nil {
return nil, fmt.Errorf("judgewalk: rank leaves: %w", err)
}
out.Usage.Add(u)
out.Requests += r
}
out.Usage.Add(u)
out.Requests += r
out.Leaves = ranked

// Gather the pages of the best leaves, in rank order, up to the
Expand Down Expand Up @@ -432,9 +563,17 @@ func (n *JudgeNavigator) Navigate(ctx context.Context, query string, leaves []Na

// More pages than the full-text budget: a coarse pass over page
// heads picks which ones deserve their whole text. Heads are
// ~150 tokens, so 120 of them is one request.
maxP := n.maxPages()
if len(pages) > maxP {
// ~150 tokens, so 120 of them is one request. Under SkimAll the heads
// are already scored, and the candidates are ordered by them.
if len(pages) > maxP && skimmed {
score := make(map[int]float64, len(heads))
for _, h := range heads {
score[h.Page.Number] = h.P
}
sort.SliceStable(pages, func(i, j int) bool { return score[pages[i].Number] > score[pages[j].Number] })
pages = pages[:maxP]
sort.Slice(pages, func(i, j int) bool { return pages[i].Number < pages[j].Number })
} else if len(pages) > maxP {
heads, u, r, err := n.rankPages(ctx, query, pages, n.headChars(), true)
if err != nil {
return nil, fmt.Errorf("judgewalk: rank page heads: %w", err)
Expand Down Expand Up @@ -514,6 +653,29 @@ func (n *JudgeNavigator) Navigate(ctx context.Context, query string, leaves []Na
return out, nil
}

// allLeafPages loads every leaf's pages once each, in document order, a
// page two leaves share kept under the first.
func allLeafPages(ctx context.Context, leaves []NavLeaf, loadPages func(ctx context.Context, leaf NavLeaf) ([]NavPage, error)) ([]NavPage, error) {
var all []NavPage
seen := map[int]bool{}
for _, l := range leaves {
ps, err := loadPages(ctx, l)
if err != nil {
return nil, fmt.Errorf("judgewalk: load %q: %w", l.Title, err)
}
for _, p := range ps {
if seen[p.Number] {
continue
}
seen[p.Number] = true
p.LeafID = l.ID
all = append(all, p)
}
}
sort.SliceStable(all, func(i, j int) bool { return all[i].Number < all[j].Number })
return all, nil
}

// countTokens estimates a string's token count for packing batches.
//
// It deliberately does NOT run the provider's tokenizer. Measured
Expand Down Expand Up @@ -691,7 +853,11 @@ func (s *JudgeWalkStrategy) selectOnPersistedPages(ctx context.Context, t *tree.
}
return ps, nil
}
nav, err := s.Navigator.Navigate(ctx, query, leaves, load)
// Every page is already in memory here, so skimming them all costs
// tokens but no loads — the case SkimAll was built for.
navigator := s.Navigator
navigator.SkimAll = true
nav, err := navigator.Navigate(ctx, query, leaves, load)
if err != nil {
return &Result{ModelUsed: "judge"}, false
}
Expand Down
Loading