package extraction import ( "context" "errors" "sort" "sync" "time" ownershipDomain "github.com/proxy-pool/proxy-pool/internal/domain/ownership" ) type Fulfillment string const ( Partial Fulfillment = "partial" AllOrNothing Fulfillment = "allOrNothing" ) type State string const ( Available State = "AVAILABLE" Extracted State = "EXTRACTED" ) var ( ErrInsufficientProxies = errors.New("insufficient proxies") ErrIdempotencyConflict = errors.New("idempotency key was reused with a different extraction request") ErrInvalidCommand = errors.New("invalid extraction command") ) type Candidate struct { ID string Protocol string Host string Port uint16 Username string Password string Region string Carrier string Upstream string OwnerWorkerID string URL string State State ExpiresAt time.Time LastCheckedAt time.Time } type Command struct { RequestID string ClientID string SourceIP string IdempotencyKey string Requested int Fulfillment Fulfillment Now time.Time MinRemainingTTL time.Duration MaxHealthCheckAge time.Duration ReserveForGateway int Protocols []string Regions []string Carriers []string Upstreams []string } type Record struct { ProxyID string ClientID string SourceIP string RequestID string Upstream string ExtractedAt time.Time ExpiresAt time.Time } type Result struct { Requested int Returned int ExtractedAt time.Time Items []Candidate } type Store interface { Extract(context.Context, Command) (Result, error) } type MemoryStore struct { mu sync.Mutex candidates map[string]Candidate records []Record idempotent map[string]idempotencyEntry nextEpoch uint64 ownership map[string]ownershipDomain.Assignment } var _ ownershipDomain.Repository = (*MemoryStore)(nil) type idempotencyEntry struct { command Command result Result } func NewMemoryStore(candidates []Candidate) *MemoryStore { items := make(map[string]Candidate, len(candidates)) for _, candidate := range candidates { items[candidate.ID] = candidate } return &MemoryStore{ candidates: items, idempotent: make(map[string]idempotencyEntry), ownership: make(map[string]ownershipDomain.Assignment), } } func (s *MemoryStore) Extract(ctx context.Context, command Command) (Result, error) { result := Result{Requested: command.Requested} if err := ctx.Err(); err != nil { return result, err } if command.Requested < 0 || command.ReserveForGateway < 0 || command.MinRemainingTTL < 0 || command.MaxHealthCheckAge < 0 || (command.Fulfillment != Partial && command.Fulfillment != AllOrNothing) { return result, ErrInvalidCommand } s.mu.Lock() defer s.mu.Unlock() if err := ctx.Err(); err != nil { return result, err } idempotencyKey := command.ClientID + "\x00" + command.IdempotencyKey if command.IdempotencyKey != "" { if committed, ok := s.idempotent[idempotencyKey]; ok { if !sameIdempotentRequest(committed.command, command) { return result, ErrIdempotencyConflict } return cloneResult(committed.result), nil } } if command.Requested <= 0 { return result, nil } eligible := make([]Candidate, 0, len(s.candidates)) for _, candidate := range s.candidates { if eligibleForExtraction(candidate, command) { eligible = append(eligible, candidate) } } sort.Slice(eligible, func(i, j int) bool { return eligible[i].ExpiresAt.After(eligible[j].ExpiresAt) }) available := len(eligible) - command.ReserveForGateway if available < 0 { available = 0 } if command.Fulfillment == AllOrNothing && available < command.Requested { return result, ErrInsufficientProxies } count := command.Requested if count > available { count = available } for i := 0; i < count; i++ { candidate := eligible[i] candidate.State = Extracted s.candidates[candidate.ID] = candidate result.Items = append(result.Items, candidate) s.records = append(s.records, Record{ ProxyID: candidate.ID, ClientID: command.ClientID, SourceIP: command.SourceIP, RequestID: command.RequestID, Upstream: candidate.Upstream, ExtractedAt: command.Now, ExpiresAt: candidate.ExpiresAt, }) } result.Returned = len(result.Items) if result.Returned > 0 { result.ExtractedAt = command.Now } if command.IdempotencyKey != "" { s.idempotent[idempotencyKey] = idempotencyEntry{ command: cloneCommand(command), result: cloneResult(result), } } return result, nil } func (s *MemoryStore) Assign(now time.Time, proxyID, workerID string, ttl time.Duration) (ownershipDomain.Assignment, error) { if proxyID == "" || workerID == "" || ttl <= 0 { return ownershipDomain.Assignment{}, ownershipDomain.ErrInvalidOwnership } s.mu.Lock() defer s.mu.Unlock() candidate, ok := s.candidates[proxyID] if current, exists := s.ownership[proxyID]; exists && current.ExpiresAt.After(now) { return ownershipDomain.Assignment{}, ownershipDomain.ErrAlreadyOwned } else if exists { if candidate.OwnerWorkerID == current.WorkerID { candidate.OwnerWorkerID = "" } delete(s.ownership, proxyID) } if !ok || candidate.State != Available || candidate.OwnerWorkerID != "" { return ownershipDomain.Assignment{}, ownershipDomain.ErrOwnershipUnavailable } s.nextEpoch++ assignment := ownershipDomain.Assignment{ ProxyID: proxyID, WorkerID: workerID, Epoch: s.nextEpoch, Version: 1, ExpiresAt: now.UTC().Add(ttl), } candidate.OwnerWorkerID = workerID s.candidates[proxyID] = candidate s.ownership[proxyID] = assignment return assignment, nil } func (s *MemoryStore) Renew(now time.Time, proxyID, workerID string, epoch uint64, ttl time.Duration) (ownershipDomain.Assignment, error) { if ttl <= 0 { return ownershipDomain.Assignment{}, ownershipDomain.ErrInvalidOwnership } s.mu.Lock() defer s.mu.Unlock() assignment, ok := s.ownership[proxyID] if !ok || assignment.WorkerID != workerID || assignment.Epoch != epoch || !assignment.ExpiresAt.After(now) { return ownershipDomain.Assignment{}, ownershipDomain.ErrStaleAssignment } assignment.ExpiresAt = now.UTC().Add(ttl) assignment.Version++ s.ownership[proxyID] = assignment return assignment, nil } func (s *MemoryStore) BeginDrain(proxyID, workerID string, epoch uint64) (ownershipDomain.Assignment, error) { s.mu.Lock() defer s.mu.Unlock() assignment, ok := s.ownership[proxyID] if !ok || assignment.WorkerID != workerID || assignment.Epoch != epoch { return ownershipDomain.Assignment{}, ownershipDomain.ErrStaleAssignment } if assignment.Draining { return assignment, nil } assignment.Draining = true assignment.Version++ s.ownership[proxyID] = assignment return assignment, nil } func (s *MemoryStore) AcknowledgeDrain(proxyID, workerID string, epoch uint64, active, reserved int64) error { if active < 0 || reserved < 0 { return ownershipDomain.ErrInvalidOwnership } s.mu.Lock() defer s.mu.Unlock() assignment, ok := s.ownership[proxyID] if !ok || assignment.WorkerID != workerID || assignment.Epoch != epoch { return ownershipDomain.ErrStaleAssignment } if !assignment.Draining { return ownershipDomain.ErrNotDraining } if active > 0 || reserved > 0 { return ownershipDomain.ErrDrainNotReady } if candidate, exists := s.candidates[proxyID]; exists && candidate.OwnerWorkerID == workerID { candidate.OwnerWorkerID = "" s.candidates[proxyID] = candidate } delete(s.ownership, proxyID) return nil } func (s *MemoryStore) Get(proxyID string) (ownershipDomain.Assignment, bool) { s.mu.Lock() defer s.mu.Unlock() assignment, ok := s.ownership[proxyID] return assignment, ok } func (s *MemoryStore) Expire(now time.Time) []ownershipDomain.Assignment { s.mu.Lock() defer s.mu.Unlock() expired := make([]ownershipDomain.Assignment, 0) for proxyID, assignment := range s.ownership { if assignment.ExpiresAt.After(now) { continue } if candidate, ok := s.candidates[proxyID]; ok && candidate.OwnerWorkerID == assignment.WorkerID { candidate.OwnerWorkerID = "" s.candidates[proxyID] = candidate } expired = append(expired, assignment) delete(s.ownership, proxyID) } sort.Slice(expired, func(i, j int) bool { return expired[i].ProxyID < expired[j].ProxyID }) return expired } func (s *MemoryStore) Records() []Record { s.mu.Lock() defer s.mu.Unlock() return append([]Record(nil), s.records...) } func cloneResult(result Result) Result { result.Items = append([]Candidate(nil), result.Items...) return result } func cloneCommand(command Command) Command { command.Protocols = append([]string(nil), command.Protocols...) command.Regions = append([]string(nil), command.Regions...) command.Carriers = append([]string(nil), command.Carriers...) command.Upstreams = append([]string(nil), command.Upstreams...) return command } func sameIdempotentRequest(left, right Command) bool { return left.Requested == right.Requested && left.Fulfillment == right.Fulfillment && equalSet(left.Protocols, right.Protocols) && equalSet(left.Regions, right.Regions) && equalSet(left.Carriers, right.Carriers) && equalSet(left.Upstreams, right.Upstreams) } func equalSet(left, right []string) bool { if len(left) != len(right) { return false } counts := make(map[string]int, len(left)) for _, value := range left { counts[value]++ } for _, value := range right { counts[value]-- if counts[value] < 0 { return false } } return true } func eligibleForExtraction(candidate Candidate, command Command) bool { if candidate.State != Available || candidate.OwnerWorkerID != "" { return false } if !candidate.ExpiresAt.IsZero() && candidate.ExpiresAt.Sub(command.Now) < command.MinRemainingTTL { return false } if command.MaxHealthCheckAge > 0 && command.Now.Sub(candidate.LastCheckedAt) > command.MaxHealthCheckAge { return false } return matches(command.Protocols, candidate.Protocol) && matches(command.Regions, candidate.Region) && matches(command.Carriers, candidate.Carrier) && matches(command.Upstreams, candidate.Upstream) } func matches(allowed []string, value string) bool { if len(allowed) == 0 { return true } for _, candidate := range allowed { if candidate == value { return true } } return false }