From 5d48690eb376f83703013bc4c59f3ee3b262ed4d Mon Sep 17 00:00:00 2001 From: Rustem Kamalov Date: Sat, 2 May 2026 16:30:56 +0300 Subject: [PATCH] feat(mega): add mode fast/any/balanced with CB latency-based fast path --- README.md | 10 +- cmd/root.go | 2 +- config.yaml | 2 +- core/circuit_breaker.go | 23 ++++ core/circuit_breaker_test.go | 7 ++ core/resilient.go | 93 +++++++++++++- core/response.go | 13 +- core/server.go | 230 +++++++++++++++++++++++++++-------- core/server_test.go | 170 ++++++++++++++++++++++++++ docs/openapi.yaml | 51 +++++++- 10 files changed, 536 insertions(+), 65 deletions(-) diff --git a/README.md b/README.md index b0e2ea9..c0a588d 100644 --- a/README.md +++ b/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" diff --git a/cmd/root.go b/cmd/root.go index 5016c92..8751984 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -16,7 +16,7 @@ import ( ) const ( - version = "0.7.6" + version = "0.7.7" defaultConfigFilename = "config" envPrefix = "OPENSERP" ) diff --git a/config.yaml b/config.yaml index 4ec9d9b..ba5dae1 100644 --- a/config.yaml +++ b/config.yaml @@ -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: diff --git a/core/circuit_breaker.go b/core/circuit_breaker.go index 3be0930..0d9f6de 100644 --- a/core/circuit_breaker.go +++ b/core/circuit_breaker.go @@ -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() diff --git a/core/circuit_breaker_test.go b/core/circuit_breaker_test.go index e676f1b..4abb7ec 100644 --- a/core/circuit_breaker_test.go +++ b/core/circuit_breaker_test.go @@ -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) diff --git a/core/resilient.go b/core/resilient.go index f23eb17..5465810 100644 --- a/core/resilient.go +++ b/core/resilient.go @@ -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) diff --git a/core/response.go b/core/response.go index 45deddb..4460525 100644 --- a/core/response.go +++ b/core/response.go @@ -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. diff --git a/core/server.go b/core/server.go index dd4c0ad..9069ff5 100644 --- a/core/server.go +++ b/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 { diff --git a/core/server_test.go b/core/server_test.go index b25498b..411eaff 100644 --- a/core/server_test.go +++ b/core/server_test.go @@ -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{ diff --git a/docs/openapi.yaml b/docs/openapi.yaml index e801e89..cb6b1c9 100644 --- a/docs/openapi.yaml +++ b/docs/openapi.yaml @@ -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]