mirror of
https://github.com/karust/openserp.git
synced 2026-08-12 20:03:29 +08:00
feat(mega): add mode fast/any/balanced with CB latency-based fast path
This commit is contained in:
10
README.md
10
README.md
@@ -73,8 +73,14 @@ Megasearch:
|
||||
# Search all configured engines
|
||||
curl "http://127.0.0.1:7000/mega/search?text=golang&limit=10"
|
||||
|
||||
# Search selected engines
|
||||
curl "http://127.0.0.1:7000/mega/search?text=golang&engines=duckduckgo,bing&limit=15"
|
||||
# Fast mode: only one fastest engine is queried
|
||||
curl "http://127.0.0.1:7000/mega/search?text=golang&mode=fast&engines=google,bing,yandex"
|
||||
|
||||
# Any mode: sequential fallback in provided order (default orded if none provided)
|
||||
curl "http://127.0.0.1:7000/mega/search?text=golang&mode=any&engines=google,yandex,bing"
|
||||
|
||||
# Balanced mode (default): parallel all engines with aggregation controls
|
||||
curl "http://127.0.0.1:7000/mega/search?text=golang&mode=balanced&dedupe=true&merge=true"
|
||||
|
||||
# Advanced filtering
|
||||
curl "http://127.0.0.1:7000/mega/search?text=golang&engines=google,bing&limit=20&date=20250101..20251231&lang=EN"
|
||||
|
||||
@@ -16,7 +16,7 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
version = "0.7.6"
|
||||
version = "0.7.7"
|
||||
defaultConfigFilename = "config"
|
||||
envPrefix = "OPENSERP"
|
||||
)
|
||||
|
||||
@@ -14,7 +14,7 @@ app:
|
||||
head: false # Headful mode
|
||||
leakless: false # Force browser process cleanup after request
|
||||
leave_head: false # Keep tabs open after request
|
||||
max_processes: 4 # Concurrent Chrome processes
|
||||
max_processes: 5 # Concurrent Chrome processes
|
||||
idle_ttl: 5m # close a Chrome that has not served traffic for this long
|
||||
mega_timeout: 90s # max total wait for /mega/* requests; slow engines return partial results
|
||||
proxies:
|
||||
|
||||
@@ -50,6 +50,8 @@ type CircuitBreaker struct {
|
||||
config CircuitBreakerConfig
|
||||
failureCount int
|
||||
successCount int
|
||||
successLatency time.Duration
|
||||
successSamples int64
|
||||
lastFailureTime time.Time
|
||||
lastStateChange time.Time
|
||||
}
|
||||
@@ -85,9 +87,18 @@ func (cb *CircuitBreaker) AllowRequest(ctx context.Context) bool {
|
||||
}
|
||||
|
||||
func (cb *CircuitBreaker) RecordSuccess(ctx context.Context) {
|
||||
cb.RecordSuccessDuration(ctx, 0)
|
||||
}
|
||||
|
||||
func (cb *CircuitBreaker) RecordSuccessDuration(ctx context.Context, elapsed time.Duration) {
|
||||
cb.mu.Lock()
|
||||
defer cb.mu.Unlock()
|
||||
|
||||
if elapsed > 0 {
|
||||
cb.successLatency += elapsed
|
||||
cb.successSamples++
|
||||
}
|
||||
|
||||
switch cb.state {
|
||||
case CircuitHalfOpen:
|
||||
cb.successCount++
|
||||
@@ -155,10 +166,22 @@ func (cb *CircuitBreaker) Stats() map[string]interface{} {
|
||||
}
|
||||
stats["retry_in"] = retryInSeconds
|
||||
}
|
||||
if cb.successSamples > 0 {
|
||||
stats["avg_response_ms"] = int64((cb.successLatency / time.Duration(cb.successSamples)) / time.Millisecond)
|
||||
}
|
||||
|
||||
return stats
|
||||
}
|
||||
|
||||
func (cb *CircuitBreaker) AvgSuccessLatency() (time.Duration, bool) {
|
||||
cb.mu.RLock()
|
||||
defer cb.mu.RUnlock()
|
||||
if cb.successSamples == 0 {
|
||||
return 0, false
|
||||
}
|
||||
return cb.successLatency / time.Duration(cb.successSamples), true
|
||||
}
|
||||
|
||||
func (cb *CircuitBreaker) setState(state CircuitState) {
|
||||
cb.state = state
|
||||
cb.lastStateChange = time.Now()
|
||||
|
||||
@@ -139,6 +139,13 @@ func TestCircuitBreaker_Stats(t *testing.T) {
|
||||
if _, ok := stats["retry_in"]; ok {
|
||||
t.Fatalf("did not expect retry_in in closed state, got: %v", stats["retry_in"])
|
||||
}
|
||||
latencyCB := NewCircuitBreaker("latency-engine", DefaultCircuitBreakerConfig())
|
||||
latencyCB.RecordSuccessDuration(context.Background(), 25*time.Millisecond)
|
||||
latencyStats := latencyCB.Stats()
|
||||
avg, ok := latencyStats["avg_response_ms"].(int64)
|
||||
if !ok || avg <= 0 {
|
||||
t.Fatalf("expected avg_response_ms int64 > 0, got: %v (%T)", latencyStats["avg_response_ms"], latencyStats["avg_response_ms"])
|
||||
}
|
||||
|
||||
openCfg := CircuitBreakerConfig{FailureThreshold: 1, RecoveryTimeout: time.Second, SuccessThreshold: 1}
|
||||
openCB := NewCircuitBreaker("open-engine", openCfg)
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/sirupsen/logrus"
|
||||
)
|
||||
@@ -169,6 +170,7 @@ func (rs *ResilientSearcher) searchWithProtection(ctx context.Context, engine Se
|
||||
policy := rs.effectivePolicyForQuery(engine.Name(), q)
|
||||
attemptMeta := rs.baseProxyMeta(policy)
|
||||
|
||||
startedAt := time.Now()
|
||||
result := RetryableSearch(ctx, rs.retryCfg, engine.Name(), func(callCtx context.Context) ([]SearchResult, error) {
|
||||
limiter := engine.GetRateLimiter()
|
||||
if limiter != nil {
|
||||
@@ -234,7 +236,7 @@ func (rs *ResilientSearcher) searchWithProtection(ctx context.Context, engine Se
|
||||
return nil, attemptMeta, result.Err
|
||||
}
|
||||
|
||||
cb.RecordSuccess(engineCtx)
|
||||
cb.RecordSuccessDuration(engineCtx, time.Since(startedAt))
|
||||
return result.Results, attemptMeta, nil
|
||||
}
|
||||
|
||||
@@ -260,6 +262,95 @@ func (rs *ResilientSearcher) searchAllImageParallelDetailed(ctx context.Context,
|
||||
return results, responded, errors
|
||||
}
|
||||
|
||||
func (rs *ResilientSearcher) searchAnyDetailed(ctx context.Context, q Query, engines []SearchEngine, isImage bool) ([]MegaSearchResult, []string, []EngineErrorDetail) {
|
||||
ctx = EnsureContext(ctx)
|
||||
|
||||
engineErrors := make([]EngineErrorDetail, 0, len(engines))
|
||||
for _, engine := range engines {
|
||||
if err := ctx.Err(); err != nil {
|
||||
engineErrors = append(engineErrors, engineErrorDetail(engine.Name(), err, q))
|
||||
break
|
||||
}
|
||||
if !engine.IsInitialized() {
|
||||
engineErrors = append(engineErrors, engineErrorDetail(engine.Name(), fmt.Errorf("not initialized"), q))
|
||||
continue
|
||||
}
|
||||
|
||||
results, _, err := rs.searchWithProtection(ctx, engine, q, isImage)
|
||||
if err != nil {
|
||||
engineErrors = append(engineErrors, engineErrorDetail(engine.Name(), err, q))
|
||||
continue
|
||||
}
|
||||
|
||||
mega := make([]MegaSearchResult, len(results))
|
||||
for i, r := range results {
|
||||
mega[i] = MegaSearchResult{SearchResult: r, Engine: engine.Name()}
|
||||
}
|
||||
return mega, []string{engine.Name()}, engineErrors
|
||||
}
|
||||
|
||||
if engineErrors == nil {
|
||||
engineErrors = []EngineErrorDetail{}
|
||||
}
|
||||
return []MegaSearchResult{}, []string{}, engineErrors
|
||||
}
|
||||
|
||||
func (rs *ResilientSearcher) searchFastestDetailed(ctx context.Context, q Query, engines []SearchEngine, isImage bool) ([]MegaSearchResult, []string, []EngineErrorDetail) {
|
||||
ctx = EnsureContext(ctx)
|
||||
|
||||
var (
|
||||
fastest SearchEngine
|
||||
bestLatency = time.Duration(1<<63 - 1)
|
||||
foundLatency bool
|
||||
unavailable []EngineErrorDetail
|
||||
candidates []SearchEngine
|
||||
)
|
||||
|
||||
for _, engine := range engines {
|
||||
cb := rs.cbManager.Get(engine.Name())
|
||||
if !engine.IsInitialized() {
|
||||
unavailable = append(unavailable, engineErrorDetail(engine.Name(), fmt.Errorf("not initialized"), q))
|
||||
continue
|
||||
}
|
||||
|
||||
engineCtx := WithEngine(ctx, engine.Name())
|
||||
if !cb.AllowRequest(engineCtx) {
|
||||
unavailable = append(unavailable, engineErrorDetail(engine.Name(), ErrCircuitOpen, q))
|
||||
continue
|
||||
}
|
||||
|
||||
candidates = append(candidates, engine)
|
||||
if avg, ok := cb.AvgSuccessLatency(); ok {
|
||||
if !foundLatency || avg < bestLatency {
|
||||
bestLatency = avg
|
||||
fastest = engine
|
||||
foundLatency = true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if len(candidates) == 0 {
|
||||
if unavailable == nil {
|
||||
unavailable = []EngineErrorDetail{}
|
||||
}
|
||||
return []MegaSearchResult{}, []string{}, unavailable
|
||||
}
|
||||
if fastest == nil {
|
||||
fastest = candidates[0]
|
||||
}
|
||||
|
||||
results, _, err := rs.searchWithProtection(ctx, fastest, q, isImage)
|
||||
if err != nil {
|
||||
return []MegaSearchResult{}, []string{}, []EngineErrorDetail{engineErrorDetail(fastest.Name(), err, q)}
|
||||
}
|
||||
|
||||
mega := make([]MegaSearchResult, len(results))
|
||||
for i, r := range results {
|
||||
mega[i] = MegaSearchResult{SearchResult: r, Engine: fastest.Name()}
|
||||
}
|
||||
return mega, []string{fastest.Name()}, []EngineErrorDetail{}
|
||||
}
|
||||
|
||||
func (rs *ResilientSearcher) runParallelDetailed(ctx context.Context, q Query, engines []SearchEngine, isImage bool) ([]MegaSearchResult, []string, []string, []EngineErrorDetail) {
|
||||
ctx = EnsureContext(ctx)
|
||||
|
||||
|
||||
@@ -11,12 +11,13 @@ type QueryEcho struct {
|
||||
|
||||
// ResponseMeta carries request-level metadata for observability and debugging.
|
||||
type ResponseMeta struct {
|
||||
RequestID string `json:"request_id"`
|
||||
RequestedAt string `json:"requested_at"`
|
||||
TookMs int64 `json:"took_ms"`
|
||||
EnginesFailed []string `json:"engines_failed"`
|
||||
EngineErrors []EngineErrorDetail `json:"engine_errors,omitempty"`
|
||||
Version string `json:"version"`
|
||||
RequestID string `json:"request_id"`
|
||||
RequestedAt string `json:"requested_at"`
|
||||
TookMs int64 `json:"took_ms"`
|
||||
EnginesResponded []string `json:"engines_responded,omitempty"`
|
||||
EnginesFailed []string `json:"engines_failed"`
|
||||
EngineErrors []EngineErrorDetail `json:"engine_errors,omitempty"`
|
||||
Version string `json:"version"`
|
||||
}
|
||||
|
||||
// EngineErrorDetail is a client-facing, sanitized per-engine failure summary.
|
||||
|
||||
230
core/server.go
230
core/server.go
@@ -804,15 +804,27 @@ type MegaSearchResult struct {
|
||||
Engine string `json:"engine"`
|
||||
}
|
||||
|
||||
const (
|
||||
megaModeBalanced = "balanced"
|
||||
megaModeAny = "any"
|
||||
megaModeFast = "fast"
|
||||
)
|
||||
|
||||
type megaRunConfig struct {
|
||||
Mode string
|
||||
Dedupe bool
|
||||
Merge bool
|
||||
}
|
||||
|
||||
func (s *Server) handleMegaSearch(c *fiber.Ctx) error {
|
||||
return s.handleMegaEndpoint(c, "search", s.resilient.searchAllParallelDetailed)
|
||||
return s.handleMegaEndpoint(c, "search")
|
||||
}
|
||||
|
||||
func (s *Server) handleMegaImage(c *fiber.Ctx) error {
|
||||
return s.handleMegaEndpoint(c, "image", s.resilient.searchAllImageParallelDetailed)
|
||||
return s.handleMegaEndpoint(c, "image")
|
||||
}
|
||||
|
||||
func (s *Server) handleMegaEndpoint(c *fiber.Ctx, action string, run func(context.Context, Query, []SearchEngine) ([]MegaSearchResult, []string, []EngineErrorDetail)) error {
|
||||
func (s *Server) handleMegaEndpoint(c *fiber.Ctx, action string) error {
|
||||
startedAt := time.Now()
|
||||
requestCtx := withRequestUsage(c.UserContext(), "mega")
|
||||
c.SetUserContext(requestCtx)
|
||||
@@ -838,8 +850,12 @@ func (s *Server) handleMegaEndpoint(c *fiber.Ctx, action string, run func(contex
|
||||
c.SetUserContext(requestCtx)
|
||||
|
||||
requestID := RequestIDFromContext(requestCtx)
|
||||
runCfg, err := parseMegaRunConfig(c)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
enginesToUse := s.resolveEngines(c.Query("engines", ""))
|
||||
enginesToUse := s.resolveEngines(requestCtx, c.Query("engines", ""))
|
||||
if len(enginesToUse) == 0 {
|
||||
return &APIError{HTTPStatus: 400, Reason: ReasonNoEngines, Message: "no valid search engines specified"}
|
||||
}
|
||||
@@ -853,22 +869,25 @@ func (s *Server) handleMegaEndpoint(c *fiber.Ctx, action string, run func(contex
|
||||
WithRequest(requestCtx).WithFields(logrus.Fields{
|
||||
"action": action,
|
||||
"engines": engineNamesJoined,
|
||||
"mode": runCfg.Mode,
|
||||
}).Debugf("Starting mega %s request for query: %s", action, q.Text)
|
||||
|
||||
cacheHitCandidates := []cacheHitCandidate{
|
||||
{
|
||||
key: s.buildMegaCacheKey(action, enginesToUse, q),
|
||||
logMessage: fmt.Sprintf("Cache hit for mega %s: engines=%s query=%s", action, engineNamesJoined, q.Text),
|
||||
},
|
||||
}
|
||||
cacheableEngines := s.megaCacheableEngines(enginesToUse)
|
||||
if len(cacheableEngines) > 0 && len(cacheableEngines) < len(enginesToUse) {
|
||||
cacheHitCandidates = append(cacheHitCandidates, cacheHitCandidate{
|
||||
key: s.buildMegaCacheKey(action, cacheableEngines, q),
|
||||
logMessage: fmt.Sprintf("Cache hit for mega %s partial set: engines=%s query=%s", action, engineNamesJoined, q.Text),
|
||||
})
|
||||
}
|
||||
if format == "json" && !ShouldBypassCacheForProxyMarket(q) {
|
||||
if format == "json" && !ShouldBypassCacheForProxyMarket(q) && runCfg.Mode != megaModeFast {
|
||||
cacheHitCandidates := []cacheHitCandidate{
|
||||
{
|
||||
key: s.buildMegaCacheKey(action, enginesToUse, q, runCfg),
|
||||
logMessage: fmt.Sprintf("Cache hit for mega %s: engines=%s query=%s mode=%s", action, engineNamesJoined, q.Text, runCfg.Mode),
|
||||
},
|
||||
}
|
||||
if runCfg.Mode == megaModeBalanced {
|
||||
cacheableEngines := s.megaCacheableEngines(enginesToUse)
|
||||
if len(cacheableEngines) > 0 && len(cacheableEngines) < len(enginesToUse) {
|
||||
cacheHitCandidates = append(cacheHitCandidates, cacheHitCandidate{
|
||||
key: s.buildMegaCacheKey(action, cacheableEngines, q, runCfg),
|
||||
logMessage: fmt.Sprintf("Cache hit for mega %s partial set: engines=%s query=%s mode=%s", action, engineNamesJoined, q.Text, runCfg.Mode),
|
||||
})
|
||||
}
|
||||
}
|
||||
if hit, err := s.tryServeCacheHit(c, startedAt, cacheHitCandidates...); hit || err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -880,9 +899,29 @@ func (s *Server) handleMegaEndpoint(c *fiber.Ctx, action string, run func(contex
|
||||
runCtx, cancel = context.WithTimeout(requestCtx, s.opts.MegaTimeout)
|
||||
defer cancel()
|
||||
}
|
||||
rawResults, _, engineErrors := run(runCtx, q, enginesToUse)
|
||||
|
||||
var (
|
||||
rawResults []MegaSearchResult
|
||||
responded []string
|
||||
engineErrors []EngineErrorDetail
|
||||
)
|
||||
switch runCfg.Mode {
|
||||
case megaModeAny:
|
||||
rawResults, responded, engineErrors = s.resilient.searchAnyDetailed(runCtx, q, enginesToUse, action == "image")
|
||||
case megaModeFast:
|
||||
rawResults, responded, engineErrors = s.resilient.searchFastestDetailed(runCtx, q, enginesToUse, action == "image")
|
||||
default:
|
||||
if action == "image" {
|
||||
rawResults, responded, engineErrors = s.resilient.searchAllImageParallelDetailed(runCtx, q, enginesToUse)
|
||||
} else {
|
||||
rawResults, responded, engineErrors = s.resilient.searchAllParallelDetailed(runCtx, q, enginesToUse)
|
||||
}
|
||||
}
|
||||
|
||||
rawResults = s.applyMegaMergePolicy(rawResults, enginesToUse, runCfg)
|
||||
|
||||
enginesFailed := engineErrorNames(engineErrors)
|
||||
if len(engineErrors) == len(enginesToUse) {
|
||||
if len(responded) == 0 {
|
||||
err := fmt.Errorf("%w: %s", ErrAllEnginesFailed, strings.Join(enginesFailed, ","))
|
||||
apiErr := searchAPIError(err, "mega", q, ProxyExecutionMeta{})
|
||||
apiErr.Meta["engine_errors"] = engineErrors
|
||||
@@ -893,17 +932,22 @@ func (s *Server) handleMegaEndpoint(c *fiber.Ctx, action string, run func(contex
|
||||
}
|
||||
|
||||
if action == "image" {
|
||||
imageResults := rawResults
|
||||
if runCfg.Dedupe {
|
||||
imageResults = s.deduplicateMegaResults(imageResults)
|
||||
}
|
||||
env := NewImageEnvelope(q, requestID, startedAt, engineNames)
|
||||
env.Meta.EnginesResponded = responded
|
||||
env.Meta.EnginesFailed = enginesFailed
|
||||
env.Meta.EngineErrors = engineErrors
|
||||
for _, r := range rawResults {
|
||||
for _, r := range imageResults {
|
||||
ectx := EnrichContext{Engine: r.Engine, Query: q}
|
||||
env.Results = append(env.Results, EnrichImageResult(r.SearchResult, ectx))
|
||||
}
|
||||
env.Finalize(startedAt, q)
|
||||
|
||||
if format == "json" && s.cache != nil {
|
||||
c.Set("X-Cache", s.cacheMegaImageResults(action, enginesToUse, q, env))
|
||||
if format == "json" && s.cache != nil && runCfg.Mode != megaModeFast {
|
||||
c.Set("X-Cache", s.cacheMegaImageResults(action, enginesToUse, q, env, runCfg))
|
||||
}
|
||||
WithRequest(requestCtx).WithFields(logrus.Fields{
|
||||
"action": action, "engines_count": len(enginesToUse), "results_count": len(env.Results),
|
||||
@@ -911,32 +955,34 @@ func (s *Server) handleMegaEndpoint(c *fiber.Ctx, action string, run func(contex
|
||||
return sendImageEnvelope(c, format, env)
|
||||
}
|
||||
|
||||
// Web search: enrich all raw results, build clusters from full set, then dedup flat list.
|
||||
allEnriched := make([]Result, 0, len(rawResults))
|
||||
for _, r := range rawResults {
|
||||
ectx := EnrichContext{Engine: r.Engine, Query: q}
|
||||
allEnriched = append(allEnriched, EnrichResult(r.SearchResult, ectx))
|
||||
webResults := rawResults
|
||||
if runCfg.Dedupe {
|
||||
webResults = s.deduplicateMegaResults(webResults)
|
||||
}
|
||||
|
||||
clusters := BuildClusters(allEnriched, len(enginesToUse))
|
||||
|
||||
// Deduplicate the flat results list by normalized URL (keep best-ranked occurrence).
|
||||
dedupedRaw := s.deduplicateMegaResults(rawResults)
|
||||
env := NewEnvelope(q, requestID, startedAt, engineNames)
|
||||
env.Meta.EnginesResponded = responded
|
||||
env.Meta.EnginesFailed = enginesFailed
|
||||
env.Meta.EngineErrors = engineErrors
|
||||
for _, r := range dedupedRaw {
|
||||
for _, r := range webResults {
|
||||
ectx := EnrichContext{Engine: r.Engine, Query: q}
|
||||
env.Results = append(env.Results, EnrichResult(r.SearchResult, ectx))
|
||||
}
|
||||
env.Finalize(startedAt, q)
|
||||
|
||||
if len(clusters) > 0 {
|
||||
env.Clusters = &clusters
|
||||
if runCfg.Merge {
|
||||
allEnriched := make([]Result, 0, len(rawResults))
|
||||
for _, r := range rawResults {
|
||||
ectx := EnrichContext{Engine: r.Engine, Query: q}
|
||||
allEnriched = append(allEnriched, EnrichResult(r.SearchResult, ectx))
|
||||
}
|
||||
clusters := BuildClusters(allEnriched, len(enginesToUse))
|
||||
if len(clusters) > 0 {
|
||||
env.Clusters = &clusters
|
||||
}
|
||||
}
|
||||
|
||||
if format == "json" && s.cache != nil {
|
||||
c.Set("X-Cache", s.cacheMegaEnvelopeResults(action, enginesToUse, q, env))
|
||||
if format == "json" && s.cache != nil && runCfg.Mode != megaModeFast {
|
||||
c.Set("X-Cache", s.cacheMegaEnvelopeResults(action, enginesToUse, q, env, runCfg))
|
||||
}
|
||||
|
||||
WithRequest(requestCtx).WithFields(logrus.Fields{
|
||||
@@ -955,6 +1001,74 @@ func engineErrorNames(details []EngineErrorDetail) []string {
|
||||
return names
|
||||
}
|
||||
|
||||
func parseMegaRunConfig(c *fiber.Ctx) (megaRunConfig, error) {
|
||||
cfg := megaRunConfig{
|
||||
Mode: megaModeBalanced,
|
||||
Dedupe: true,
|
||||
Merge: true,
|
||||
}
|
||||
|
||||
mode := strings.ToLower(strings.TrimSpace(c.Query("mode", megaModeBalanced)))
|
||||
switch mode {
|
||||
case "", megaModeBalanced:
|
||||
cfg.Mode = megaModeBalanced
|
||||
case megaModeAny:
|
||||
cfg.Mode = megaModeAny
|
||||
case megaModeFast:
|
||||
cfg.Mode = megaModeFast
|
||||
default:
|
||||
return megaRunConfig{}, errInvalidParam("mode: must be one of fast, any, balanced")
|
||||
}
|
||||
|
||||
if raw := strings.TrimSpace(c.Query("dedupe", "")); raw != "" {
|
||||
value, err := strconv.ParseBool(raw)
|
||||
if err != nil {
|
||||
return megaRunConfig{}, errInvalidParam(fmt.Sprintf("dedupe: %v", err))
|
||||
}
|
||||
cfg.Dedupe = value
|
||||
}
|
||||
|
||||
if raw := strings.TrimSpace(c.Query("merge", "")); raw != "" {
|
||||
value, err := strconv.ParseBool(raw)
|
||||
if err != nil {
|
||||
return megaRunConfig{}, errInvalidParam(fmt.Sprintf("merge: %v", err))
|
||||
}
|
||||
cfg.Merge = value
|
||||
}
|
||||
|
||||
return cfg, nil
|
||||
}
|
||||
|
||||
func (s *Server) applyMegaMergePolicy(results []MegaSearchResult, engines []SearchEngine, cfg megaRunConfig) []MegaSearchResult {
|
||||
if cfg.Merge || len(results) == 0 {
|
||||
return results
|
||||
}
|
||||
|
||||
byEngine := make(map[string]struct{}, len(results))
|
||||
for _, r := range results {
|
||||
byEngine[r.Engine] = struct{}{}
|
||||
}
|
||||
|
||||
selected := ""
|
||||
for _, eng := range engines {
|
||||
if _, ok := byEngine[eng.Name()]; ok {
|
||||
selected = eng.Name()
|
||||
break
|
||||
}
|
||||
}
|
||||
if selected == "" {
|
||||
return []MegaSearchResult{}
|
||||
}
|
||||
|
||||
filtered := make([]MegaSearchResult, 0, len(results))
|
||||
for _, r := range results {
|
||||
if r.Engine == selected {
|
||||
filtered = append(filtered, r)
|
||||
}
|
||||
}
|
||||
return filtered
|
||||
}
|
||||
|
||||
func (s *Server) handleListEngines(c *fiber.Ctx) error {
|
||||
var engines []map[string]interface{}
|
||||
|
||||
@@ -981,7 +1095,7 @@ func (s *Server) handleListEngines(c *fiber.Ctx) error {
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Server) resolveEngines(enginesParam string) []SearchEngine {
|
||||
func (s *Server) resolveEngines(ctx context.Context, enginesParam string) []SearchEngine {
|
||||
if enginesParam == "" {
|
||||
return s.searchEngines
|
||||
}
|
||||
@@ -989,22 +1103,37 @@ func (s *Server) resolveEngines(enginesParam string) []SearchEngine {
|
||||
var enginesToUse []SearchEngine
|
||||
seen := make(map[string]bool)
|
||||
engineNames := strings.Split(enginesParam, ",")
|
||||
for _, engineName := range engineNames {
|
||||
engineName = strings.TrimSpace(strings.ToLower(engineName))
|
||||
if engineName == "" || seen[engineName] {
|
||||
for _, name := range engineNames {
|
||||
name = strings.TrimSpace(strings.ToLower(name))
|
||||
if name == "" || seen[name] {
|
||||
continue
|
||||
}
|
||||
name = resolveEngineAlias(name)
|
||||
matched := false
|
||||
for _, engine := range s.searchEngines {
|
||||
if strings.ToLower(engine.Name()) == engineName {
|
||||
if strings.ToLower(engine.Name()) == name {
|
||||
enginesToUse = append(enginesToUse, engine)
|
||||
seen[engineName] = true
|
||||
seen[name] = true
|
||||
matched = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !matched {
|
||||
WithRequest(ctx).Warnf("Unknown engine %q requested, skipping", name)
|
||||
}
|
||||
}
|
||||
return enginesToUse
|
||||
}
|
||||
|
||||
func resolveEngineAlias(name string) string {
|
||||
switch name {
|
||||
case "duck", "ddg":
|
||||
return "duckduckgo"
|
||||
default:
|
||||
return name
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Server) deduplicateMegaResults(results []MegaSearchResult) []MegaSearchResult {
|
||||
urlMap := make(map[string]MegaSearchResult)
|
||||
order := []string{}
|
||||
@@ -1112,7 +1241,7 @@ func (s *Server) cacheJSON(cacheKey string, payload interface{}) bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func (s *Server) cacheMegaEnvelopeResults(action string, enginesToUse []SearchEngine, q Query, env *Envelope) string {
|
||||
func (s *Server) cacheMegaEnvelopeResults(action string, enginesToUse []SearchEngine, q Query, env *Envelope, cfg megaRunConfig) string {
|
||||
if s.cache == nil {
|
||||
return ""
|
||||
}
|
||||
@@ -1129,13 +1258,13 @@ func (s *Server) cacheMegaEnvelopeResults(action string, enginesToUse []SearchEn
|
||||
s.cache.RecordBypass()
|
||||
return "BYPASS"
|
||||
}
|
||||
if s.cacheJSON(s.buildMegaCacheKey(action, cacheEngines, q), env) {
|
||||
if s.cacheJSON(s.buildMegaCacheKey(action, cacheEngines, q, cfg), env) {
|
||||
return "MISS"
|
||||
}
|
||||
return "BYPASS"
|
||||
}
|
||||
|
||||
func (s *Server) cacheMegaImageResults(action string, enginesToUse []SearchEngine, q Query, env *ImageEnvelope) string {
|
||||
func (s *Server) cacheMegaImageResults(action string, enginesToUse []SearchEngine, q Query, env *ImageEnvelope, cfg megaRunConfig) string {
|
||||
if s.cache == nil {
|
||||
return ""
|
||||
}
|
||||
@@ -1152,13 +1281,13 @@ func (s *Server) cacheMegaImageResults(action string, enginesToUse []SearchEngin
|
||||
s.cache.RecordBypass()
|
||||
return "BYPASS"
|
||||
}
|
||||
if s.cacheJSON(s.buildMegaCacheKey(action, cacheEngines, q), env) {
|
||||
if s.cacheJSON(s.buildMegaCacheKey(action, cacheEngines, q, cfg), env) {
|
||||
return "MISS"
|
||||
}
|
||||
return "BYPASS"
|
||||
}
|
||||
|
||||
func (s *Server) buildMegaCacheKey(action string, engines []SearchEngine, q Query) string {
|
||||
func (s *Server) buildMegaCacheKey(action string, engines []SearchEngine, q Query, cfg megaRunConfig) string {
|
||||
names := make([]string, 0, len(engines))
|
||||
for _, eng := range engines {
|
||||
names = append(names, strings.ToLower(strings.TrimSpace(eng.Name())))
|
||||
@@ -1177,7 +1306,8 @@ func (s *Server) buildMegaCacheKey(action string, engines []SearchEngine, q Quer
|
||||
last = name
|
||||
}
|
||||
|
||||
return BuildCacheKey("mega:"+strings.Join(uniq, ","), action, q)
|
||||
prefix := fmt.Sprintf("mega:%s:%t:%t:%s", cfg.Mode, cfg.Merge, cfg.Dedupe, strings.Join(uniq, ","))
|
||||
return BuildCacheKey(prefix, action, q)
|
||||
}
|
||||
|
||||
func (s *Server) megaCacheableEngines(engines []SearchEngine) []SearchEngine {
|
||||
|
||||
@@ -970,6 +970,176 @@ func TestMegaSearchDeduplicatesRepeatedEngineNames(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestMegaSearchInvalidModeReturnsBadRequest(t *testing.T) {
|
||||
engine := &engineMock{name: "google", initialized: true}
|
||||
opts := DefaultServerOptions()
|
||||
opts.Resilience.Retry.MaxRetries = 0
|
||||
srv := NewServerWithOptions("127.0.0.1", 7181, opts, engine)
|
||||
|
||||
resp := request(t, srv, "/mega/search?text=golang&mode=turbo")
|
||||
if resp.StatusCode != http.StatusBadRequest {
|
||||
t.Fatalf("expected 400 for invalid mode, got %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
var payload JSONErrorResponse
|
||||
if err := json.NewDecoder(resp.Body).Decode(&payload); err != nil {
|
||||
t.Fatalf("decode error payload: %v", err)
|
||||
}
|
||||
if payload.Reason != ReasonInvalidParam {
|
||||
t.Fatalf("expected reason=%q, got %q", ReasonInvalidParam, payload.Reason)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMegaSearchAnyModeStopsAfterFirstSuccessInRequestedOrder(t *testing.T) {
|
||||
google := &engineMock{
|
||||
name: "google",
|
||||
initialized: true,
|
||||
searchFn: func(_ context.Context, q Query) ([]SearchResult, error) {
|
||||
return nil, errors.New("google failed")
|
||||
},
|
||||
}
|
||||
yandex := &engineMock{
|
||||
name: "yandex",
|
||||
initialized: true,
|
||||
searchFn: func(_ context.Context, q Query) ([]SearchResult, error) {
|
||||
return []SearchResult{{Rank: 1, URL: "https://example.com/yandex", Title: "yandex"}}, nil
|
||||
},
|
||||
}
|
||||
bing := &engineMock{
|
||||
name: "bing",
|
||||
initialized: true,
|
||||
searchFn: func(_ context.Context, q Query) ([]SearchResult, error) {
|
||||
return []SearchResult{{Rank: 1, URL: "https://example.com/bing", Title: "bing"}}, nil
|
||||
},
|
||||
}
|
||||
|
||||
opts := DefaultServerOptions()
|
||||
opts.Resilience.Retry.MaxRetries = 0
|
||||
opts.CacheTTL = 0
|
||||
srv := NewServerWithOptions("127.0.0.1", 7182, opts, google, yandex, bing)
|
||||
|
||||
resp := request(t, srv, "/mega/search?text=golang&mode=any&engines=google,yandex,bing")
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
t.Fatalf("expected 200, got %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
var env Envelope
|
||||
if err := json.NewDecoder(resp.Body).Decode(&env); err != nil {
|
||||
t.Fatalf("decode envelope: %v", err)
|
||||
}
|
||||
if len(env.Results) == 0 || env.Results[0].Engine != "yandex" {
|
||||
t.Fatalf("expected yandex result in any mode, got %#v", env.Results)
|
||||
}
|
||||
|
||||
if google.searchCalls != 1 || yandex.searchCalls != 1 || bing.searchCalls != 0 {
|
||||
t.Fatalf("expected google=1 yandex=1 bing=0 calls, got google=%d yandex=%d bing=%d", google.searchCalls, yandex.searchCalls, bing.searchCalls)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMegaSearchFastModeUsesFastestEngineFromCircuitStats(t *testing.T) {
|
||||
google := &engineMock{
|
||||
name: "google",
|
||||
initialized: true,
|
||||
searchFn: func(_ context.Context, q Query) ([]SearchResult, error) {
|
||||
time.Sleep(25 * time.Millisecond)
|
||||
return []SearchResult{{Rank: 1, URL: "https://example.com/google", Title: "google"}}, nil
|
||||
},
|
||||
}
|
||||
yandex := &engineMock{
|
||||
name: "yandex",
|
||||
initialized: true,
|
||||
searchFn: func(_ context.Context, q Query) ([]SearchResult, error) {
|
||||
time.Sleep(1 * time.Millisecond)
|
||||
return []SearchResult{{Rank: 1, URL: "https://example.com/yandex", Title: "yandex"}}, nil
|
||||
},
|
||||
}
|
||||
|
||||
opts := DefaultServerOptions()
|
||||
opts.Resilience.Retry.MaxRetries = 0
|
||||
opts.CacheTTL = 0
|
||||
srv := NewServerWithOptions("127.0.0.1", 7183, opts, google, yandex)
|
||||
|
||||
_ = request(t, srv, "/google/search?text=warm-google")
|
||||
_ = request(t, srv, "/yandex/search?text=warm-yandex")
|
||||
|
||||
google.mu.Lock()
|
||||
beforeGoogle := google.searchCalls
|
||||
google.mu.Unlock()
|
||||
yandex.mu.Lock()
|
||||
beforeYandex := yandex.searchCalls
|
||||
yandex.mu.Unlock()
|
||||
|
||||
resp := request(t, srv, "/mega/search?text=golang&mode=fast&engines=google,yandex")
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
t.Fatalf("expected 200, got %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
var env Envelope
|
||||
if err := json.NewDecoder(resp.Body).Decode(&env); err != nil {
|
||||
t.Fatalf("decode envelope: %v", err)
|
||||
}
|
||||
if len(env.Results) == 0 || env.Results[0].Engine != "yandex" {
|
||||
t.Fatalf("expected fast mode to use yandex, got %#v", env.Results)
|
||||
}
|
||||
|
||||
google.mu.Lock()
|
||||
afterGoogle := google.searchCalls
|
||||
google.mu.Unlock()
|
||||
yandex.mu.Lock()
|
||||
afterYandex := yandex.searchCalls
|
||||
yandex.mu.Unlock()
|
||||
|
||||
if afterGoogle-beforeGoogle != 0 || afterYandex-beforeYandex != 1 {
|
||||
t.Fatalf("expected delta calls google=0 yandex=1, got google=%d yandex=%d", afterGoogle-beforeGoogle, afterYandex-beforeYandex)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMegaSearchBalancedModeMergeAndDedupeFlags(t *testing.T) {
|
||||
google := &engineMock{
|
||||
name: "google",
|
||||
initialized: true,
|
||||
searchFn: func(_ context.Context, q Query) ([]SearchResult, error) {
|
||||
return []SearchResult{{Rank: 1, URL: "https://example.com/shared", Title: "google"}}, nil
|
||||
},
|
||||
}
|
||||
bing := &engineMock{
|
||||
name: "bing",
|
||||
initialized: true,
|
||||
searchFn: func(_ context.Context, q Query) ([]SearchResult, error) {
|
||||
return []SearchResult{{Rank: 2, URL: "https://example.com/shared", Title: "bing"}}, nil
|
||||
},
|
||||
}
|
||||
|
||||
opts := DefaultServerOptions()
|
||||
opts.Resilience.Retry.MaxRetries = 0
|
||||
opts.CacheTTL = 0
|
||||
srv := NewServerWithOptions("127.0.0.1", 7184, opts, google, bing)
|
||||
|
||||
respNoDedupe := request(t, srv, "/mega/search?text=golang&mode=balanced&dedupe=false&engines=google,bing")
|
||||
if respNoDedupe.StatusCode != http.StatusOK {
|
||||
t.Fatalf("expected 200 for dedupe=false, got %d", respNoDedupe.StatusCode)
|
||||
}
|
||||
var envNoDedupe Envelope
|
||||
if err := json.NewDecoder(respNoDedupe.Body).Decode(&envNoDedupe); err != nil {
|
||||
t.Fatalf("decode dedupe=false envelope: %v", err)
|
||||
}
|
||||
if len(envNoDedupe.Results) != 2 {
|
||||
t.Fatalf("expected 2 results with dedupe=false, got %d", len(envNoDedupe.Results))
|
||||
}
|
||||
|
||||
respNoMerge := request(t, srv, "/mega/search?text=golang&mode=balanced&merge=false&engines=google,bing")
|
||||
if respNoMerge.StatusCode != http.StatusOK {
|
||||
t.Fatalf("expected 200 for merge=false, got %d", respNoMerge.StatusCode)
|
||||
}
|
||||
var envNoMerge Envelope
|
||||
if err := json.NewDecoder(respNoMerge.Body).Decode(&envNoMerge); err != nil {
|
||||
t.Fatalf("decode merge=false envelope: %v", err)
|
||||
}
|
||||
if len(envNoMerge.Results) != 1 || envNoMerge.Results[0].Engine != "google" {
|
||||
t.Fatalf("expected merge=false to keep first requested engine results, got %#v", envNoMerge.Results)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMegaSearchCachesForHealthySubsetWhenOneCircuitIsOpen(t *testing.T) {
|
||||
google := &engineMock{name: "google", initialized: true}
|
||||
bing := &engineMock{
|
||||
|
||||
@@ -217,10 +217,13 @@ paths:
|
||||
get:
|
||||
tags: [Mega]
|
||||
operationId: megaSearch
|
||||
summary: Search across multiple engines in parallel
|
||||
summary: Search across multiple engines with selectable execution mode
|
||||
description: >
|
||||
Results are deduplicated by normalized URL. The `clusters` field groups
|
||||
results that appeared in multiple engines, scored by cross-engine agreement.
|
||||
Mode controls engine execution strategy: `balanced` (default) queries all
|
||||
selected engines in parallel, `any` runs engines sequentially in requested
|
||||
order until first success, and `fast` queries only the fastest engine based
|
||||
on circuit-breaker average response time stats.
|
||||
In `balanced` mode, `dedupe` and `merge` tune aggregation behavior.
|
||||
Partial failures are surfaced in `meta.engines_failed` and
|
||||
`meta.engine_errors`. If all selected engines fail, the endpoint returns
|
||||
a 502 with per-engine error details. Use `?format=markdown|text|ndjson`
|
||||
@@ -236,6 +239,9 @@ paths:
|
||||
- $ref: "#/components/parameters/FilterQuery"
|
||||
- $ref: "#/components/parameters/AnswersQuery"
|
||||
- $ref: "#/components/parameters/EnginesQuery"
|
||||
- $ref: "#/components/parameters/MegaModeQuery"
|
||||
- $ref: "#/components/parameters/MegaDedupeQuery"
|
||||
- $ref: "#/components/parameters/MegaMergeQuery"
|
||||
- $ref: "#/components/parameters/FormatQuery"
|
||||
- $ref: "#/components/parameters/UseProxyHeader"
|
||||
- $ref: "#/components/parameters/ProxyURLHeader"
|
||||
@@ -293,7 +299,7 @@ paths:
|
||||
get:
|
||||
tags: [Mega]
|
||||
operationId: megaImageSearch
|
||||
summary: Image search across multiple engines in parallel
|
||||
summary: Image search across multiple engines with selectable execution mode
|
||||
parameters:
|
||||
- $ref: "#/components/parameters/TextQuery"
|
||||
- $ref: "#/components/parameters/LangQuery"
|
||||
@@ -305,6 +311,9 @@ paths:
|
||||
- $ref: "#/components/parameters/FilterQuery"
|
||||
- $ref: "#/components/parameters/AnswersQuery"
|
||||
- $ref: "#/components/parameters/EnginesQuery"
|
||||
- $ref: "#/components/parameters/MegaModeQuery"
|
||||
- $ref: "#/components/parameters/MegaDedupeQuery"
|
||||
- $ref: "#/components/parameters/MegaMergeQuery"
|
||||
- $ref: "#/components/parameters/FormatQuery"
|
||||
- $ref: "#/components/parameters/UseProxyHeader"
|
||||
- $ref: "#/components/parameters/ProxyURLHeader"
|
||||
@@ -567,6 +576,37 @@ components:
|
||||
schema:
|
||||
type: string
|
||||
example: google,bing,duckduckgo
|
||||
MegaModeQuery:
|
||||
name: mode
|
||||
in: query
|
||||
required: false
|
||||
description: >
|
||||
Mega execution mode. `balanced` (default) runs all selected engines in parallel.
|
||||
`any` runs selected engines sequentially in request order until first success.
|
||||
`fast` runs only one engine: the fastest by circuit-breaker average response time.
|
||||
schema:
|
||||
type: string
|
||||
enum: [balanced, any, fast]
|
||||
default: balanced
|
||||
MegaDedupeQuery:
|
||||
name: dedupe
|
||||
in: query
|
||||
required: false
|
||||
description: >
|
||||
Enable deduplication by normalized URL. Default `true`.
|
||||
schema:
|
||||
type: boolean
|
||||
default: true
|
||||
MegaMergeQuery:
|
||||
name: merge
|
||||
in: query
|
||||
required: false
|
||||
description: >
|
||||
Merge results from all successful engines into one flat list. Default `true`.
|
||||
When `false`, only the first requested engine that returned results is kept.
|
||||
schema:
|
||||
type: boolean
|
||||
default: true
|
||||
FormatQuery:
|
||||
name: format
|
||||
in: query
|
||||
@@ -1507,6 +1547,9 @@ components:
|
||||
retry_in:
|
||||
type: integer
|
||||
description: Seconds until next half-open attempt (present when state is open).
|
||||
avg_response_ms:
|
||||
type: integer
|
||||
description: Average successful engine response time in milliseconds.
|
||||
CircuitBreakerStatsResponse:
|
||||
type: object
|
||||
required: [circuit_breakers]
|
||||
|
||||
Reference in New Issue
Block a user