Files
openserp/core/proxy.go
Rustem Kamalov 9cc69da758 feat(cli): clean search, add extract/format flags; harden engines & proxy rotation; fix bugs; update docs
- Add structured `search [engine] [query]` CLI: --limit/--lang/--region/--site/--file, --format (json|text|markdown|ndjson), --extract N, --search-timeout; Envelope and route logs to stderr with a --quiet default (fixes stdout pollution)
- Unify engines behind a single engineSpec registry (CLI + serve share it)
- Unify the extract knob to bool-or-int `extract=N` (drop extract_top); CLI and HTTP share core batch extraction, raw/rendered fetch, and clamp helpers
- Engines: Ecosia CF captcha detection (raw + browser), Yandex progressive-result wait, Google PAA poll + Has() existence probes, Bing title/desc attribute fallbacks
- Proxy: rotate challenged proxies out of the tag pool for one retry (X-Proxy-Attempts); browser health-ping skip window; opt-in WaitStable
2026-06-16 03:51:37 +03:00

701 lines
18 KiB
Go

package core
import (
"context"
"errors"
"fmt"
"net/url"
"sort"
"strings"
"sync"
"time"
"github.com/sirupsen/logrus"
)
const (
ProxyRuntimeBrowser = "browser"
ProxyRuntimeRaw = "raw"
ProxyModeOff = "off"
ProxyModeTagPool = "tag_pool"
ProxyModeRequestURL = "request_url"
DefaultProxyFailureThreshold = 3
ProxyOverrideDirect = "direct"
// ProxyPoolQuarantineDuration is how long an exhausted tag pool stays quarantined
// before a single probe proxy is re-enabled for recovery testing.
ProxyPoolQuarantineDuration = 5 * time.Minute
// ProxyChallengeCooldown is how long a captcha/blocked proxy is deprioritized
// in rotation. It does not degrade health, so the proxy is still served if
// it's the only one left.
ProxyChallengeCooldown = 2 * time.Minute
)
var supportedProxySchemes = map[string]struct{}{
"http": {},
"https": {},
"socks5": {},
"socks5h": {},
}
var ErrProxyUnavailable = errors.New("proxy unavailable")
type ProxyPolicy struct {
Mode string `json:"mode" mapstructure:"mode"`
Tag string `json:"tag,omitempty" mapstructure:"tag"`
}
type ProxyEntryConfig struct {
URL string `json:"url" mapstructure:"url"`
Tags []string `json:"tags" mapstructure:"tags"`
}
type ProxiesHealthConfig struct {
FailureThreshold int `json:"failure_threshold" mapstructure:"failure_threshold"`
}
type ProxiesConfig struct {
Global string `json:"global,omitempty" mapstructure:"global"`
Entries []ProxyEntryConfig `json:"entries" mapstructure:"entries"`
Health ProxiesHealthConfig `json:"health" mapstructure:"health"`
AllowRequestProxyURL bool `json:"allow_request_proxy_url" mapstructure:"allow_request_proxy_url"`
Lanes ProxyLanesConfig `json:"lanes" mapstructure:"lanes"`
}
type ProxyConfig struct {
Runtime string // raw or browser runtime behavior
Proxies ProxiesConfig // canonical proxy inventory
EnginePolicies map[string]string // engine-specific proxy tags
Registry *ProxyRegistry // optional shared registry from caller
}
type ProxyTagSummary struct {
Configured int `json:"configured"`
Healthy int `json:"healthy"`
}
type ProxyStatsEntry struct {
Proxy string `json:"proxy"`
Tags []string `json:"tags"`
Healthy bool `json:"healthy"`
Failures int `json:"failures"`
Disabled bool `json:"disabled"`
}
type ProxyEngineStats struct {
Tag string `json:"tag,omitempty"`
SelectedProxy string `json:"selected_proxy"`
}
type ProxyStats struct {
ConfiguredCount int `json:"configured_count"`
HealthyCount int `json:"healthy_count"`
UnhealthyCount int `json:"unhealthy_count"`
RequestProxyURLEnabled bool `json:"request_proxy_url_enabled"`
Lanes LaneStats `json:"lanes"`
BrowserProcesses BrowserPoolStats `json:"browser_processes"`
Tags map[string]ProxyTagSummary `json:"tags"`
Entries []ProxyStatsEntry `json:"entries"`
Engines map[string]ProxyEngineStats `json:"engines,omitempty"`
}
type proxyState struct {
url string
tags []string
failures int
disabled bool
// challengedUntil deprioritizes (but does not disable) this proxy in
// rotation after a captcha/block. See ReportChallenged.
challengedUntil time.Time
}
type ProxyRegistry struct {
mu sync.Mutex
states map[string]*proxyState
order []string
tagIndex map[string][]string
nextByTag map[string]int
tagQuarantine map[string]time.Time // tag to earliest time pool may probe again
failureThreshold int
}
func DefaultProxiesConfig() ProxiesConfig {
return ProxiesConfig{
Global: "",
Entries: []ProxyEntryConfig{},
Health: ProxiesHealthConfig{FailureThreshold: DefaultProxyFailureThreshold},
Lanes: DefaultProxyLanesConfig(),
}
}
func DefaultProxyConfig() ProxyConfig {
return ProxyConfig{
Runtime: ProxyRuntimeBrowser,
Proxies: DefaultProxiesConfig(),
EnginePolicies: map[string]string{},
}
}
func NormalizeProxyConfig(cfg ProxyConfig) (ProxyConfig, error) {
cfg.Runtime = normalizeProxyRuntime(cfg.Runtime)
var err error
cfg.Proxies, err = NormalizeProxiesConfig(cfg.Proxies)
if err != nil {
return cfg, err
}
if cfg.EnginePolicies == nil {
cfg.EnginePolicies = map[string]string{}
}
normalizedEnginePolicies := make(map[string]string, len(cfg.EnginePolicies))
for rawEngine, rawTag := range cfg.EnginePolicies {
engine := normalizeEngineName(rawEngine)
if engine == "" {
continue
}
tag := normalizeTag(rawTag)
if tag == "" {
continue
}
normalizedEnginePolicies[engine] = tag
}
cfg.EnginePolicies = normalizedEnginePolicies
if cfg.Registry == nil {
if len(cfg.Proxies.Entries) > 0 {
registry, err := NewProxyRegistry(cfg.Proxies.Entries, cfg.Proxies.Health.FailureThreshold)
if err != nil {
return cfg, err
}
cfg.Registry = registry
}
}
return cfg, nil
}
func NormalizeProxiesConfig(cfg ProxiesConfig) (ProxiesConfig, error) {
global, err := NormalizeProxyURL(cfg.Global)
if err != nil {
return cfg, fmt.Errorf("invalid proxies.global: %w", err)
}
cfg.Global = global
failureThreshold := cfg.Health.FailureThreshold
if failureThreshold <= 0 {
failureThreshold = DefaultProxyFailureThreshold
}
normalizedEntries := make([]ProxyEntryConfig, 0, len(cfg.Entries))
entryByURL := make(map[string]int, len(cfg.Entries))
for i, rawEntry := range cfg.Entries {
proxyURL, err := NormalizeProxyURL(rawEntry.URL)
if err != nil {
return cfg, fmt.Errorf("invalid proxies.entries[%d].url: %w", i, err)
}
if proxyURL == "" {
return cfg, fmt.Errorf("invalid proxies.entries[%d].url: value is required", i)
}
tags, err := normalizeProxyTags(rawEntry.Tags)
if err != nil {
return cfg, fmt.Errorf("invalid proxies.entries[%d].tags: %w", i, err)
}
if idx, ok := entryByURL[proxyURL]; ok {
normalizedEntries[idx].Tags = mergeTags(normalizedEntries[idx].Tags, tags)
continue
}
normalizedEntries = append(normalizedEntries, ProxyEntryConfig{
URL: proxyURL,
Tags: tags,
})
entryByURL[proxyURL] = len(normalizedEntries) - 1
}
cfg.Entries = normalizedEntries
cfg.Health = ProxiesHealthConfig{FailureThreshold: failureThreshold}
cfg.Lanes = NormalizeProxyLanesConfig(cfg.Lanes)
return cfg, nil
}
func NormalizeProxyURL(raw string) (string, error) {
raw = strings.TrimSpace(raw)
if raw == "" {
return "", nil
}
parsed, err := url.Parse(raw)
if err != nil {
return "", err
}
if parsed.Scheme == "" {
return "", fmt.Errorf("proxy URL must include a scheme")
}
if parsed.Host == "" {
return "", fmt.Errorf("proxy URL must include a host")
}
parsed.Scheme = strings.ToLower(parsed.Scheme)
if _, ok := supportedProxySchemes[parsed.Scheme]; !ok {
return "", fmt.Errorf("unsupported proxy scheme %q", parsed.Scheme)
}
return parsed.String(), nil
}
func NormalizeProxyURLs(rawURLs []string) ([]string, error) {
normalized := make([]string, 0, len(rawURLs))
seen := make(map[string]struct{}, len(rawURLs))
for _, raw := range rawURLs {
proxyURL, err := NormalizeProxyURL(raw)
if err != nil {
return nil, err
}
if proxyURL == "" {
continue
}
if _, ok := seen[proxyURL]; ok {
continue
}
seen[proxyURL] = struct{}{}
normalized = append(normalized, proxyURL)
}
return normalized, nil
}
func MaskProxyURL(raw string) string {
proxyURL, err := NormalizeProxyURL(raw)
if err != nil || proxyURL == "" {
return "invalid-proxy"
}
parsed, err := url.Parse(proxyURL)
if err != nil {
return "invalid-proxy"
}
return fmt.Sprintf("%s://%s", parsed.Scheme, parsed.Host)
}
func ResolveEffectiveProxyPolicy(globalProxyURL string, engineTag string) ProxyPolicy {
if strings.TrimSpace(globalProxyURL) != "" {
return ProxyPolicy{Mode: ProxyModeTagPool}
}
tag := normalizeTag(engineTag)
if tag == "" {
return ProxyPolicy{Mode: ProxyModeOff}
}
return ProxyPolicy{Mode: ProxyModeTagPool, Tag: tag}
}
func NormalizeProxyTag(raw string) (string, error) {
tag := normalizeTag(raw)
if tag == "" {
return "", fmt.Errorf("value is required")
}
return tag, nil
}
func NormalizeProxyRequestOverride(raw string) (string, error) {
override := normalizeTag(raw)
if override == "" {
return "", nil
}
if override == ProxyOverrideDirect {
return ProxyOverrideDirect, nil
}
return NormalizeProxyTag(override)
}
func IsAuthenticatedSocksProxyURL(raw string) bool {
normalized, err := NormalizeProxyURL(raw)
if err != nil || normalized == "" {
return false
}
parsed, err := url.Parse(normalized)
if err != nil {
return false
}
if (parsed.Scheme == "socks5" || parsed.Scheme == "socks5h") && parsed.User != nil {
return true
}
return false
}
func NewProxyRegistry(entries []ProxyEntryConfig, failureThreshold int) (*ProxyRegistry, error) {
if failureThreshold <= 0 {
failureThreshold = DefaultProxyFailureThreshold
}
states := make(map[string]*proxyState, len(entries))
order := make([]string, 0, len(entries))
tagIndex := make(map[string][]string)
for idx, entry := range entries {
proxyURL, err := NormalizeProxyURL(entry.URL)
if err != nil {
return nil, fmt.Errorf("invalid proxy registry entry[%d] url: %w", idx, err)
}
if proxyURL == "" {
return nil, fmt.Errorf("invalid proxy registry entry[%d] url: value is required", idx)
}
tags, err := normalizeProxyTags(entry.Tags)
if err != nil {
return nil, fmt.Errorf("invalid proxy registry entry[%d] tags: %w", idx, err)
}
states[proxyURL] = &proxyState{url: proxyURL, tags: tags}
order = append(order, proxyURL)
for _, tag := range tags {
tagIndex[tag] = append(tagIndex[tag], proxyURL)
}
}
return &ProxyRegistry{
states: states,
order: order,
tagIndex: tagIndex,
nextByTag: make(map[string]int, len(tagIndex)),
tagQuarantine: make(map[string]time.Time, len(tagIndex)),
failureThreshold: failureThreshold,
}, nil
}
func (r *ProxyRegistry) NextByTag(tag string) string {
return r.NextByTagWithContext(context.TODO(), tag)
}
func (r *ProxyRegistry) NextByTagWithContext(ctx context.Context, tag string) string {
tag = normalizeTag(tag)
if tag == "" {
return ""
}
r.mu.Lock()
defer r.mu.Unlock()
urls := r.tagIndex[tag]
if len(urls) == 0 {
return ""
}
if r.allDisabledLocked(urls) {
now := time.Now()
quarantineUntil := r.tagQuarantine[tag]
if quarantineUntil.IsZero() {
r.startTagQuarantineLocked(ctx, tag, now)
return ""
}
if now.Before(quarantineUntil) {
// Pool is still in quarantine; refuse to serve any proxy.
WithRequest(ctx).WithField("proxy_tag", tag).WithField("quarantine_until", quarantineUntil.Format(time.RFC3339)).
Warn("Proxy tag pool in quarantine, no proxy served")
return ""
}
delete(r.tagQuarantine, tag)
// Quarantine elapsed: re-enable one proxy as a recovery probe.
probe := r.leastFailedLocked(urls)
if probe != "" {
state := r.states[probe]
state.disabled = false
state.failures = 0
WithRequest(ctx).WithFields(logrus.Fields{
"proxy_tag": tag,
"proxy": MaskProxyURL(probe),
}).Warn("Proxy tag pool quarantine elapsed, probing one proxy for recovery")
}
}
now := time.Now()
start := r.nextByTag[tag]
// First pass skips challenged proxies; second pass relaxes that so a
// challenged-but-healthy proxy is still served rather than failing.
for _, skipChallenged := range []bool{true, false} {
for i := 0; i < len(urls); i++ {
idx := (start + i) % len(urls)
proxyURL := urls[idx]
state := r.states[proxyURL]
if state.disabled {
continue
}
if skipChallenged && now.Before(state.challengedUntil) {
continue
}
r.nextByTag[tag] = (idx + 1) % len(urls)
WithRequest(ctx).WithFields(logrus.Fields{
"proxy_tag": tag,
"proxy": MaskProxyURL(proxyURL),
}).Debugf("Selected proxy for tag=%s: %s", tag, MaskProxyURL(proxyURL))
return proxyURL
}
}
// All proxies are disabled and no probe could be selected.
r.startTagQuarantineLocked(ctx, tag, time.Now())
return ""
}
// ReportFailure increments the failure counter for proxyURL. The proxy is
// disabled once the failure threshold is reached. If the owning tag pool
// becomes fully exhausted, a quarantine timer is started so that
// NextByTagWithContext will not immediately re-enable all proxies.
//
// Only proxy-network errors (ErrProxyConnect, ErrProxyAuth, ErrTimeout) should
// degrade proxy health. Callers must not call this for captcha or parser errors.
func (r *ProxyRegistry) ReportFailure(ctx context.Context, proxyURL string) {
proxyURL, err := NormalizeProxyURL(proxyURL)
if err != nil || proxyURL == "" {
return
}
r.mu.Lock()
defer r.mu.Unlock()
state, ok := r.states[proxyURL]
if !ok {
return
}
state.failures++
if state.failures >= r.failureThreshold {
state.disabled = true
WithRequest(ctx).WithFields(logrus.Fields{
"failure_count": state.failures,
"proxy": MaskProxyURL(proxyURL),
}).Warnf("Disabled proxy after %d failures: %s", state.failures, MaskProxyURL(proxyURL))
// If all proxies in every shared tag are now disabled, start quarantine.
now := time.Now()
for _, tag := range state.tags {
if r.allDisabledLocked(r.tagIndex[tag]) {
quarantineUntil := r.tagQuarantine[tag]
if quarantineUntil.IsZero() || !now.Before(quarantineUntil) {
r.startTagQuarantineLocked(ctx, tag, now)
}
}
}
}
}
func (r *ProxyRegistry) ReportSuccess(_ context.Context, proxyURL string) {
proxyURL, err := NormalizeProxyURL(proxyURL)
if err != nil || proxyURL == "" {
return
}
r.mu.Lock()
defer r.mu.Unlock()
state, ok := r.states[proxyURL]
if !ok {
return
}
state.failures = 0
state.disabled = false
// Clear quarantine for any tag this proxy belongs to; at least one proxy is healthy again.
for _, tag := range state.tags {
delete(r.tagQuarantine, tag)
}
}
// ReportChallenged deprioritizes a captcha/blocked proxy for
// ProxyChallengeCooldown without degrading its health (unlike ReportFailure, it
// never disables the proxy or trips quarantine), so the next attempt prefers a
// different IP.
func (r *ProxyRegistry) ReportChallenged(ctx context.Context, proxyURL string) {
proxyURL, err := NormalizeProxyURL(proxyURL)
if err != nil || proxyURL == "" {
return
}
r.mu.Lock()
defer r.mu.Unlock()
state, ok := r.states[proxyURL]
if !ok {
return
}
state.challengedUntil = time.Now().Add(ProxyChallengeCooldown)
WithRequest(ctx).WithField("proxy", MaskProxyURL(proxyURL)).
Debugf("Deprioritized challenged proxy for %s: %s", ProxyChallengeCooldown, MaskProxyURL(proxyURL))
}
func (r *ProxyRegistry) HasHealthyProxyForTag(tag string) bool {
return r.HealthyCountForTag(tag) > 0
}
// HealthyCountForTag returns how many non-disabled proxies the tag pool holds
// (a challenged proxy still counts — it's usable, just deprioritized).
func (r *ProxyRegistry) HealthyCountForTag(tag string) int {
tag = normalizeTag(tag)
if tag == "" {
return 0
}
r.mu.Lock()
defer r.mu.Unlock()
count := 0
for _, proxyURL := range r.tagIndex[tag] {
if state, ok := r.states[proxyURL]; ok && !state.disabled {
count++
}
}
return count
}
func (r *ProxyRegistry) BuildStats() ProxyStats {
r.mu.Lock()
defer r.mu.Unlock()
stats := ProxyStats{
Tags: map[string]ProxyTagSummary{},
Entries: make([]ProxyStatsEntry, 0, len(r.order)),
}
for _, proxyURL := range r.order {
state := r.states[proxyURL]
healthy := !state.disabled
if healthy {
stats.HealthyCount++
} else {
stats.UnhealthyCount++
}
stats.ConfiguredCount++
stats.Entries = append(stats.Entries, ProxyStatsEntry{
Proxy: MaskProxyURL(state.url),
Tags: append([]string(nil), state.tags...),
Healthy: healthy,
Failures: state.failures,
Disabled: state.disabled,
})
for _, tag := range state.tags {
summary := stats.Tags[tag]
summary.Configured++
if healthy {
summary.Healthy++
}
stats.Tags[tag] = summary
}
}
return stats
}
func (r *ProxyRegistry) allDisabledLocked(urls []string) bool {
if len(urls) == 0 {
return false
}
for _, proxyURL := range urls {
if state, ok := r.states[proxyURL]; ok && !state.disabled {
return false
}
}
return true
}
func (r *ProxyRegistry) startTagQuarantineLocked(ctx context.Context, tag string, now time.Time) {
quarantineUntil := now.Add(ProxyPoolQuarantineDuration)
r.tagQuarantine[tag] = quarantineUntil
WithRequest(ctx).WithFields(logrus.Fields{
"proxy_tag": tag,
"quarantine_until": quarantineUntil.Format(time.RFC3339),
}).Warnf("Proxy tag pool %q fully exhausted, quarantined for %s", tag, ProxyPoolQuarantineDuration)
}
// leastFailedLocked returns the URL of the disabled proxy with the lowest
// failure count, which is the cheapest probe candidate. Must hold r.mu.
func (r *ProxyRegistry) leastFailedLocked(urls []string) string {
best := ""
bestFailures := -1
for _, proxyURL := range urls {
state, ok := r.states[proxyURL]
if !ok {
continue
}
if bestFailures < 0 || state.failures < bestFailures {
best = proxyURL
bestFailures = state.failures
}
}
return best
}
func normalizeProxyRuntime(runtime string) string {
switch strings.ToLower(strings.TrimSpace(runtime)) {
case ProxyRuntimeRaw:
return ProxyRuntimeRaw
default:
return ProxyRuntimeBrowser
}
}
func normalizeProxyTags(tags []string) ([]string, error) {
if len(tags) == 0 {
return nil, fmt.Errorf("at least one tag is required")
}
seen := make(map[string]struct{}, len(tags))
normalized := make([]string, 0, len(tags))
for _, rawTag := range tags {
tag := normalizeTag(rawTag)
if tag == "" {
continue
}
if _, ok := seen[tag]; ok {
continue
}
seen[tag] = struct{}{}
normalized = append(normalized, tag)
}
if len(normalized) == 0 {
return nil, fmt.Errorf("at least one non-empty tag is required")
}
sort.Strings(normalized)
return normalized, nil
}
func normalizeTag(raw string) string {
return strings.ToLower(strings.TrimSpace(raw))
}
func mergeTags(base []string, additional []string) []string {
combined := make(map[string]struct{}, len(base)+len(additional))
for _, tag := range base {
combined[tag] = struct{}{}
}
for _, tag := range additional {
combined[tag] = struct{}{}
}
merged := make([]string, 0, len(combined))
for tag := range combined {
merged = append(merged, tag)
}
sort.Strings(merged)
return merged
}
func normalizeEngineName(raw string) string {
return strings.ToLower(strings.TrimSpace(raw))
}