fix: parallelize Workbench session summary reads (#1900)

This commit is contained in:
Lyon
2026-06-22 16:05:32 +08:00
committed by GitHub
parent f4d961a33e
commit 153727d111
+29 -9
View File
@@ -1166,19 +1166,39 @@ func (s *Server) querySummaryFacts(ctx context.Context, sessions []map[string]an
sessionIDs := uniqueStrings(mapSessions(sessions, sessionID))
traceIDs := uniqueStrings(mapSessions(sessions, lastTraceID))
facts := emptyFacts()
type summaryFactsResult struct {
kind string
messageSummaries []any
facts map[string][]any
err error
}
resultCount := 0
results := make(chan summaryFactsResult, 2)
if len(sessionIDs) > 0 {
messageSummaries, err := s.queryMessageSummaries(ctx, sessionIDs)
if err != nil {
return nil, err
}
facts["messageSummaries"] = messageSummaries
resultCount++
go func() {
messageSummaries, err := s.queryMessageSummaries(ctx, sessionIDs)
results <- summaryFactsResult{kind: "messageSummaries", messageSummaries: messageSummaries, err: err}
}()
}
if len(traceIDs) > 0 {
byTrace, err := s.queryFacts(ctx, factQuery{TraceIDs: traceIDs, Families: []string{"turns", "checkpoints"}, Limit: len(traceIDs), TurnsOrder: "updated_desc", CheckpointsOrder: "updated_desc"})
if err != nil {
return nil, err
resultCount++
go func() {
byTrace, err := s.queryFacts(ctx, factQuery{TraceIDs: traceIDs, Families: []string{"turns", "checkpoints"}, Limit: len(traceIDs), TurnsOrder: "updated_desc", CheckpointsOrder: "updated_desc"})
results <- summaryFactsResult{kind: "traceFacts", facts: byTrace, err: err}
}()
}
for i := 0; i < resultCount; i++ {
result := <-results
if result.err != nil {
return nil, result.err
}
switch result.kind {
case "messageSummaries":
facts["messageSummaries"] = result.messageSummaries
case "traceFacts":
facts = mergeFacts(facts, result.facts)
}
facts = mergeFacts(facts, byTrace)
}
return facts, nil
}