-
Notifications
You must be signed in to change notification settings - Fork 63
Expand file tree
/
Copy pathmain.py
More file actions
1686 lines (1488 loc) · 63.6 KB
/
Copy pathmain.py
File metadata and controls
1686 lines (1488 loc) · 63.6 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
#!/usr/bin/env python3
# encoding: utf-8
"""
ProxyPool - SOCKS5 代理池一键工具
子命令:
fetch 从公开源抓取 SOCKS5 代理,合并去重写入 socks.txt
check 多线程检测 socks.txt 中代理存活性,写入 alive.txt
serve 启动 HTTP 代理服务(上游随机选 SOCKS5),供 Burp 等使用
run 一键全流程:fetch -> check -> serve
用法示例:
python3 main.py fetch --proxy socks5://127.0.0.1:1080
python3 main.py check --threads 100
python3 main.py serve --port 8082
python3 main.py run --proxy socks5://127.0.0.1:1080
"""
import argparse
import base64
import json
import os
import random
import re
import select
import signal
import socket
import threading
import time
import urllib.request
from collections import deque
from queue import Queue, Empty
from urllib.error import URLError
# ============================================================
# 共享配置
# ============================================================
# fetch 抓取源:均实测可达,返回纯文本 ip:port 列表
SOURCES = {
"TheSpeedX": "https://raw.githubusercontent.com/TheSpeedX/SOCKS-List/master/socks5.txt",
"Monosans": "https://raw.githubusercontent.com/monosans/proxy-list/main/proxies/socks5.txt",
"ProxyScrape": "https://api.proxyscrape.com/v2/?request=displayproxies&protocol=socks5"
"&timeout=10000&ssl=all&anonymity=all",
}
DEFAULT_SOCKS_FILE = "socks.txt"
DEFAULT_ALIVE_FILE = "alive.txt"
DEFAULT_REGION_FILE = "region.json"
CONFIG_FILE = "config.json"
# 配置默认值(启动/UI 未指定时使用)
CONFIG_DEFAULTS = {
"listen": "127.0.0.1",
"port": 8082, # 代理端口(HTTP+SOCKS5 入站)
"stats_port": 8083, # 控制台端口
"upstream_type": "socks5",
"max_clients": 100,
"fail_threshold": 3,
"timeout": 6,
"retries": 3,
"threads": 100, # check 线程数
"fetch_proxy": "", # fetch 时经由的代理(留空直连)
"socks_file": DEFAULT_SOCKS_FILE,
"alive_file": DEFAULT_ALIVE_FILE,
"region_file": DEFAULT_REGION_FILE,
}
# 改动后需要重启监听 socket 的字段
RESTART_KEYS = ("listen", "port")
def _count_lines(path):
"""安全读取文件行数(去空行)。文件不存在返回 0。"""
try:
with open(path, "r", encoding="utf-8", errors="ignore") as f:
return sum(1 for l in f if l.strip())
except OSError:
return 0
class Config:
"""
运行配置:持久化到 config.json,支持 UI/命令行覆盖。
线程安全:所有读写加锁。
"""
def __init__(self, path=CONFIG_FILE):
self.path = path
self.lock = threading.RLock()
self.data = dict(CONFIG_DEFAULTS)
self.load()
def load(self):
"""从 config.json 读取,与默认值合并。"""
try:
with open(self.path, "r", encoding="utf-8") as f:
saved = json.load(f)
if isinstance(saved, dict):
with self.lock:
for k in CONFIG_DEFAULTS:
if k in saved:
self.data[k] = saved[k]
except (FileNotFoundError, ValueError):
pass # 首次运行或文件损坏,用默认值
def save(self):
"""写回 config.json。"""
with self.lock:
data = dict(self.data)
try:
with open(self.path, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2)
except OSError as e:
log.warn("保存 config.json 失败: {}".format(e))
def get(self, key, default=None):
with self.lock:
return self.data.get(key, default)
def get_all(self):
with self.lock:
return dict(self.data)
def update(self, changes):
"""
批量更新字段。返回 set(被改动的字段名)。
调用方据返回值判断是否需要重启监听。
"""
changed = set()
with self.lock:
for k, v in changes.items():
if k in CONFIG_DEFAULTS and self.data.get(k) != v:
self.data[k] = v
changed.add(k)
if changed:
self.save()
return changed
def apply_overrides(self, overrides):
"""启动时用命令行参数覆盖(仅覆盖非 None 的项),不保存。"""
with self.lock:
for k, v in overrides.items():
if v is not None and k in CONFIG_DEFAULTS:
self.data[k] = v
# IP:Port 合法性校验
PROXY_RE = re.compile(r"^(\d{1,3})\.(\d{1,3})\.(\d{1,3})\.(\d{1,3}):(\d{1,5})$")
UA = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36"
def banner(title):
print("\n" + "=" * 50)
print(" " + title)
print("=" * 50)
# 代理模式
MODES = ("auto", "sticky", "manual")
class Logger:
"""
统一日志:同时输出到控制台与环形缓冲。
缓冲供 Web 控制台 /logs 拉取。每条日志带递增序号,支持 ?since 增量查询。
"""
def __init__(self, maxlen=500):
self._buf = deque(maxlen=maxlen)
self._lock = threading.Lock()
self._seq = 0
def log(self, msg, level="INFO"):
line = "[{}] {}".format(level, msg)
with self._lock:
self._seq += 1
self._buf.append((self._seq, level, msg))
# 控制台输出(不影响缓冲)
print(line)
def info(self, msg):
self.log(msg, "INFO")
def warn(self, msg):
self.log(msg, "WARN")
def error(self, msg):
self.log(msg, "ERROR")
def snapshot(self, since=0):
"""返回序号 > since 的日志条目列表:[{seq, level, msg}]。"""
with self._lock:
return [{"seq": s, "level": lv, "msg": m}
for s, lv, m in self._buf if s > since]
# 全局日志实例(serve 及其工作线程共享)
log = Logger()
# ============================================================
# fetch 子命令:抓取聚合
# ============================================================
def _valid_proxy(line):
"""校验一行是否为合法 ip:port(含数值范围检查)。"""
m = PROXY_RE.match(line.strip())
if not m:
return None
octets = [int(x) for x in m.groups()[:4]]
port = int(m.group(5))
if any(o > 255 for o in octets) or not (1 <= port <= 65535):
return None
return "{}.{}.{}.{}:{}".format(*octets, port)
def _setup_urllib_proxy(proxy_url, logger=None):
"""
如指定 --proxy,用 PySocks 把 urllib.request 走 SOCKS5。
成功返回 True;无代理或失败返回 False(直连)。
"""
lg = logger or log
if not proxy_url:
return False
# 仅支持 socks5/socks4/http
try:
import socks # PySocks
except ImportError:
lg.warn("未安装 PySocks,无法使用 --proxy,改为直连")
return False
m = re.match(r"^(socks5|socks5h|socks4|http)://([^:]+):(\d+)$", proxy_url)
if not m:
lg.warn("--proxy 格式错误,应为 socks5://host:port,改为直连")
return False
scheme, host, port = m.group(1), m.group(2), int(m.group(3))
# 全局劫持 socket 的方式仅作用于本进程抓取阶段,且为单线程顺序抓取,可接受
type_map = {"socks5": socks.PROXY_TYPE_SOCKS5, "socks5h": socks.PROXY_TYPE_SOCKS5,
"socks4": socks.PROXY_TYPE_SOCKS4, "http": socks.PROXY_TYPE_HTTP}
socks.set_default_proxy(type_map[scheme], host, port)
socket.socket = socks.socksocket
lg.info("抓取阶段经由代理 {}://{}:{}".format(scheme, host, port))
return True
def fetch_from_source(name, url, timeout, logger=None):
"""抓单个源,返回合法代理列表。失败返回空列表。"""
try:
req = urllib.request.Request(url, headers={"User-Agent": UA})
with urllib.request.urlopen(req, timeout=timeout) as resp:
raw = resp.read().decode("utf-8", "ignore")
except (URLError, socket.timeout, OSError) as e:
if logger:
logger.warn("{} 抓取失败: {}".format(name, e))
return []
found = []
for line in raw.splitlines():
v = _valid_proxy(line)
if v:
found.append(v)
return found
def run_fetch(proxy=None, output=DEFAULT_SOCKS_FILE, timeout=20, overwrite=False,
logger=None):
"""
fetch 核心逻辑(供 CLI 与控制台共享)。
返回 (总数, 新增数)。所有进度通过 logger 反馈。
"""
lg = logger or log
_setup_urllib_proxy(proxy, lg)
existing = set()
if os.path.exists(output) and not overwrite:
with open(output, "r", encoding="utf-8", errors="ignore") as f:
existing = {l.strip() for l in f if _valid_proxy(l)}
lg.info("已有 {} 条: {}".format(len(existing), output))
all_proxies = set(existing)
for name, url in SOURCES.items():
lg.info("抓取源: {}".format(name))
got = fetch_from_source(name, url, timeout, lg)
got = list(set(got)) # 源内去重
lg.info("获取 {} 条 (去重后) from {}".format(len(got), name))
all_proxies |= set(got)
ordered = sorted(all_proxies)
with open(output, "w", encoding="utf-8") as f:
f.write("\n".join(ordered) + "\n")
new = len(ordered) - len(existing)
lg.info("写入 {} 共 {} 条 (新增 {} 条)".format(output, len(ordered), new))
return len(ordered), new
def cmd_fetch(args):
banner("FETCH - 抓取 SOCKS5 代理")
run_fetch(proxy=args.proxy, output=args.output, timeout=args.timeout,
overwrite=args.overwrite)
# ============================================================
# 地区查询(检测存活时顺带获取出口 IP 地区)
# ============================================================
# 经代理访问,拿出口 IP 的地区(这些站点返回的是"代理出口 IP"的地区)
REGION_URLS = [
"http://www.cip.cc/", # 返回 IP\t: x.x.x.x / 地址\t: 中国 xx xx
"http://myip.ipip.net/", # 返回 当前 IP:x 来自于:中国 xx xx
]
# 国家判定:地区字符串含这些关键词视为中国
CN_KEYWORDS = ("中国", "China", "CN", "中华")
def _parse_region(text):
"""
从 cip.cc / myip.ipip.net 的响应文本提取地区字符串。
返回地区字符串(如 '中国 广西 南宁')或 None。
严格过滤:结果不得含 HTML 标签、引号、尖括号等,长度合理。
"""
if not text:
return None
# 若响应明显是 HTML 页面(含 <html/<head/<title 等),视为非预期格式
if re.search(r"<\s*(html|head|title|script|meta)\b", text, re.I):
return None
def _clean(s):
"""清洗候选地区:去标签、去非法字符、校验长度。"""
s = re.sub(r"<[^>]+>", "", s) # 去任意 HTML 标签
s = s.strip()
# 合法地区只含中文、字母、空格、连字符;不含引号/尖括号/斜杠
if not s:
return None
if re.search(r'[<>"\'=/\\]', s):
return None
if len(s) > 40 or len(s) < 2:
return None
return s
# cip.cc 格式: "地址\t: 中国 广西 南宁"
m = re.search(r"地址\s*:?\s*([^\n\r]+)", text)
if m:
region = _clean(m.group(1))
if region:
return region
# myip.ipip.net 格式: "来自于:中国 广西 南宁 移动"
m = re.search(r"来自于[::]\s*([^\n\r]+)", text)
if m:
parts = _clean(m.group(1))
if parts:
# 去掉末尾运营商词,保留到市级
words = parts.split()
return " ".join(words[:3]) if words else None
return None
def _country_of(region):
"""根据地区字符串判定国家代码:中国→CN,否则→OTHER。"""
if not region:
return "UNKNOWN"
for kw in CN_KEYWORDS:
if kw in region:
return "CN"
return "OTHER"
class RegionStore:
"""
地区元数据存储:ip:port -> {region, country}。
线程安全,持久化到 region.json。alive.txt 保持纯 ip:port 不变。
"""
def __init__(self, path=DEFAULT_REGION_FILE):
self.path = path
self.lock = threading.Lock()
self.data = {}
self.load()
def load(self):
try:
with open(self.path, "r", encoding="utf-8") as f:
d = json.load(f)
if isinstance(d, dict):
self.data = d
except (FileNotFoundError, ValueError):
pass
def save(self):
try:
with open(self.path, "w", encoding="utf-8") as f:
json.dump(self.data, f, ensure_ascii=False, indent=2)
except OSError as e:
log.warn("保存 region.json 失败: {}".format(e))
def update(self, proxy, region):
"""更新单个代理地区。region 为 None 时记录为未知。"""
with self.lock:
self.data[proxy] = {"region": region or "未知", "country": _country_of(region)}
def batch_save(self):
with self.lock:
self.save()
def get(self, proxy):
with self.lock:
return self.data.get(proxy)
def all(self):
with self.lock:
return dict(self.data)
def remove(self, proxy):
with self.lock:
self.data.pop(proxy, None)
def count_by_country(self):
"""返回 {CN: n, OTHER: n, UNKNOWN: n} 统计。"""
with self.lock:
cnt = {"CN": 0, "OTHER": 0, "UNKNOWN": 0}
for v in self.data.values():
c = v.get("country", "UNKNOWN")
cnt[c] = cnt.get(c, 0) + 1
return cnt
# ============================================================
# check 子命令:存活检测
# ============================================================
class ProxyChecker(threading.Thread):
"""多线程检测代理存活性,顺带通过代理获取出口地区。"""
def __init__(self, check_queue, alive_list, lock, timeout, counters,
region_store, logger=None):
threading.Thread.__init__(self)
self.check_queue = check_queue
self.alive_list = alive_list
self.lock = lock
self.timeout = timeout
self.counters = counters # {"done":0, "total":N} 共享计数
self.region_store = region_store
self.logger = logger
def run(self):
while True:
try:
target = self.check_queue.get_nowait()
except Empty:
return
try:
self.check_one(target)
finally:
self.check_queue.task_done()
with self.lock:
self.counters["done"] += 1
def check_one(self, proxy):
proxies = {"http": "socks5://" + proxy, "https": "socks5://" + proxy}
headers = {"User-Agent": UA, "Connection": "close"}
try:
import requests
requests.packages.urllib3.disable_warnings()
region = None
# 优先通过代理访问地区站点(一举两得:验证存活 + 拿出口地区)
for url in REGION_URLS:
try:
r = requests.get(url, headers=headers, proxies=proxies,
timeout=self.timeout, verify=False)
if r.status_code == 200:
region = _parse_region(r.text)
if region:
break
except Exception:
continue
# 若地区站点都失败,回退用 baidu 验证存活(地区未知)
if region is None:
try:
r = requests.get("http://www.baidu.com", headers=headers,
proxies=proxies, timeout=self.timeout, verify=False)
if r.status_code != 200:
return # 不存活
except Exception:
return # 不存活
# 存活:记录 + 存地区
with self.lock:
self.alive_list.append(proxy)
self.region_store.update(proxy, region)
country = _country_of(region)
tag = region if region else "未知地区"
if self.logger:
self.logger.info("存活: {} [{}]".format(proxy, tag))
except Exception:
pass
def run_check(input_file=DEFAULT_SOCKS_FILE, output=DEFAULT_ALIVE_FILE,
threads=100, timeout=3, logger=None, region_path=DEFAULT_REGION_FILE):
"""
check 核心逻辑(供 CLI 与控制台共享)。
检测存活时顺带通过代理获取出口地区,存入 region_path。
返回 (存活数, 总数)。进度通过 logger 和返回值反馈。
"""
lg = logger or log
try:
import requests # noqa
except ImportError:
lg.error("需要 requests 库: pip install requests")
return 0, 0
# 读取并去重输入
try:
with open(input_file, "r", encoding="utf-8", errors="ignore") as f:
raw = [l.strip() for l in f if l.strip()]
except FileNotFoundError:
lg.error("输入文件不存在: {}".format(input_file))
return 0, 0
uniq = list(set(raw))
with open(input_file, "w", encoding="utf-8") as f:
f.write("\n".join(uniq) + "\n")
queue = Queue()
for p in uniq:
queue.put(p)
total = len(uniq)
lg.info("待检测 {} 条,线程数 {}(同时获取出口地区)".format(total, threads))
alive_list = []
lock = threading.Lock()
counters = {"done": 0, "total": total}
region_store = RegionStore(region_path)
n_threads = min(threads, total) or 1
pool_check = [ProxyChecker(queue, alive_list, lock, timeout, counters,
region_store, lg) for _ in range(n_threads)]
for t in pool_check:
t.start()
for t in pool_check:
t.join()
with open(output, "w", encoding="utf-8") as f:
f.write("\n".join(alive_list) + "\n")
# 清理已不存活的地区记录,保存
alive_set = set(alive_list)
with region_store.lock:
region_store.data = {k: v for k, v in region_store.data.items()
if k in alive_set}
region_store.save()
# 地区统计
cc = region_store.count_by_country()
lg.info("存活 {} / {} 条,已写入 {}(中国 {},其他 {},未知 {})".format(
len(alive_list), total, output, cc["CN"], cc["OTHER"], cc["UNKNOWN"]))
return len(alive_list), total
def cmd_check(args):
banner("CHECK - 检测代理存活性")
run_check(input_file=args.input, output=args.output,
threads=args.threads, timeout=args.timeout)
# ============================================================
# serve 子命令:HTTP 代理服务
# ============================================================
class ProxyPool:
"""
线程安全的代理池:内存缓存、随机选取、运行时熔断、文件热加载、统计。
熔断策略:单代理连续失败 fail_threshold 次即移出池子;成功一次即清零。
热加载:pick() 前检查 alive.txt 的 mtime,变化则自动重新载入。
"""
def __init__(self, path, fail_threshold=3, region_store=None, socks_file=None):
self.path = path
self.socks_file = socks_file or DEFAULT_SOCKS_FILE
self.fail_threshold = fail_threshold
self.region_store = region_store # RegionStore 实例(可为 None)
self.lock = threading.RLock()
self.proxies = [] # 当前可用 ip:port 列表
self.fail_count = {} # ip:port -> 连续失败次数
self.removed = set() # 本轮被熔断移除的代理(避免热加载后立刻加回)
self._mtime = 0 # 上次载入时的文件 mtime
# 模式状态
self.mode = "auto" # auto / sticky / manual
self.current = None # 当前正在使用 / sticky 粘住的代理
self.locked = None # manual 模式锁定的代理
# 地区过滤:None=不过滤, "CN"=仅中国, "OTHER"=仅非中国
self.region_filter = None
# 统计
self.stat_requests = 0 # 总 pick 次数
self.stat_success = 0 # 上游连接成功次数
self.stat_fail = 0 # 上游连接失败次数
self.stat_circuit_open = 0 # 触发熔断的代理数
self.stat_reload = 0 # 热加载次数
self.stat_rotate = 0 # 手动轮换次数
self.load()
def load(self):
"""从文件载入代理列表,保留当前熔断状态。"""
try:
st = os.stat(self.path)
self._mtime = st.st_mtime
with open(self.path, "r", encoding="utf-8", errors="ignore") as f:
new = [l.strip() for l in f if _valid_proxy(l)]
except FileNotFoundError:
new = []
with self.lock:
# 合并:新列表中、且未被熔断的,纳入池子
self.proxies = [p for p in new if p not in self.removed]
# 清理已不在文件的失败计数
self.fail_count = {p: c for p, c in self.fail_count.items() if p in new}
# 清理已不在池中的 current/locked
if self.current and self.current not in self.proxies:
self.current = None
if self.locked and self.locked not in self.proxies:
self.locked = None
if not self.proxies:
log.warn("代理池为空 ({} 不存在或无有效行)".format(self.path))
else:
log.info("代理池载入 {} 条: {}".format(len(self.proxies), self.path))
def _maybe_reload(self):
"""mtime 变化则重新载入。持锁调用前不锁,内部加锁。"""
try:
mtime = os.stat(self.path).st_mtime
except OSError:
return
if mtime != self._mtime:
log.info("检测到 {} 变更,热加载...".format(self.path))
self.stat_reload += 1
self.load()
def _eligible(self, candidates):
"""
按地区过滤候选列表。region_filter 为 None 时不过滤。
无 region_store 或无记录的代理,过滤时不予剔除(宽松保留)。
"""
rf = self.region_filter
if not rf or not self.region_store:
return candidates
result = []
for p in candidates:
info = self.region_store.get(p)
if info is None:
result.append(p) # 无地区记录,保留
elif info.get("country") == rf:
result.append(p)
return result
def pick(self):
"""
按 mode 返回一个代理;池空返回 None。触发热加载检查。
- auto: 每次随机(受地区过滤约束)
- sticky: 优先复用 current(成功粘住的代理),失效则随机重选
- manual: 固定返回 locked,被熔断则返回 None(不自动换)
"""
self._maybe_reload()
with self.lock:
self.stat_requests += 1
if not self.proxies:
return None
if self.mode == "manual":
if self.locked and self.locked in self.proxies:
return self.locked
return None # 锁定的被熔断,不自动换
if self.mode == "sticky":
if self.current and self.current in self.proxies:
return self.current
# current 失效,从符合地区的候选中随机重选并粘住
cands = self._eligible(self.proxies) or self.proxies
self.current = random.choice(cands)
return self.current
# auto:从符合地区的候选中随机
cands = self._eligible(self.proxies) or self.proxies
return random.choice(cands)
def mark_success(self, proxy):
"""标记代理连接成功,清零失败计数;sticky 模式下粘住它。"""
with self.lock:
self.stat_success += 1
self.fail_count.pop(proxy, None)
if self.mode == "sticky":
self.current = proxy
def mark_fail(self, proxy):
"""
标记代理连接失败。连续失败达阈值则移出池子(熔断)。
返回 True 表示本次触发了熔断。
"""
with self.lock:
self.stat_fail += 1
c = self.fail_count.get(proxy, 0) + 1
self.fail_count[proxy] = c
if c >= self.fail_threshold and proxy in self.proxies:
try:
self.proxies.remove(proxy)
except ValueError:
pass
self.removed.add(proxy)
self.stat_circuit_open += 1
# 被熔断的若是 current/locked,清掉
if self.current == proxy:
self.current = None
if self.locked == proxy:
self.locked = None
log.warn("[⚡熔断] {} 连续失败 {} 次,移出池子 (剩余 {})".format(
proxy, c, len(self.proxies)))
return True
return False
def set_mode(self, mode):
"""切换模式。切到 manual 时自动锁定当前 current(若有)。"""
if mode not in MODES:
return False
with self.lock:
self.mode = mode
if mode == "manual":
# 锁定当前使用中的代理(优先 current)
self.locked = self.current
else:
self.locked = None
log.info("模式切换为: {}".format(mode))
return True
def set_region_filter(self, country):
"""设置地区过滤:None=全部, 'CN'=仅中国, 'OTHER'=仅非中国。"""
if country not in (None, "CN", "OTHER"):
return False
with self.lock:
self.region_filter = country
# 切换过滤后清掉可能不符的 current
if country and self.current:
info = self.region_store.get(self.current) if self.region_store else None
if info and info.get("country") != country:
self.current = None
label = {None: "全部", "CN": "仅中国", "OTHER": "仅其他"}.get(country, str(country))
log.info("地区过滤: {}".format(label))
return True
def rotate(self):
"""
手动轮换到下一个随机代理(受地区过滤约束)。
更新 current;manual 模式下同时更新 locked。返回新代理或 None。
"""
self._maybe_reload()
with self.lock:
if not self.proxies:
return None
cands = self._eligible(self.proxies) or self.proxies
# 尽量换一个不同的;只有 1 个候选时只能返回同一个
if len(cands) > 1:
choices = [p for p in cands if p != self.current]
new = random.choice(choices) if choices else random.choice(cands)
else:
new = cands[0]
self.current = new
if self.mode == "manual":
self.locked = new
self.stat_rotate += 1
log.info("[🔄轮换] 切换到 {}".format(new))
return new
def current_ip(self):
"""返回当前生效代理:manual 用 locked,否则用 current。"""
with self.lock:
if self.mode == "manual":
return self.locked
return self.current
def stats(self):
"""返回统计快照(字典),含模式、当前代理、地区过滤与统计。"""
with self.lock:
cur = self.locked if self.mode == "manual" else self.current
d = {
"pool_size": len(self.proxies), # 可用:存活且未熔断
"removed": len(self.removed), # 已熔断
"requests": self.stat_requests,
"success": self.stat_success,
"fail": self.stat_fail,
"circuit_open": self.stat_circuit_open,
"reload": self.stat_reload,
"rotate": self.stat_rotate,
"mode": self.mode,
"current": cur,
"region_filter": self.region_filter,
}
# 抓取总数 / 存活总数(读文件,不随熔断变)
d["socks_total"] = _count_lines(self.socks_file)
d["alive_total"] = _count_lines(self.path)
# 地区统计(来自 region_store,不加 pool 锁)
if self.region_store:
d["region_stats"] = self.region_store.count_by_country()
return d
def reset_circuits(self):
"""清除所有熔断状态,把移除的代理加回池子。"""
with self.lock:
self.proxies.extend(self.removed)
self.removed.clear()
self.fail_count.clear()
n = len(self.proxies)
log.info("熔断状态已重置,池子恢复至 {} 条".format(n))
return n
class Header:
"""读取并解析客户端请求头。first 为协议识别阶段已读出的首字节(可省)。"""
def __init__(self, conn, first=b""):
self._method = None
header = first
try:
while True:
data = conn.recv(4096)
header = b"%s%s" % (header, data)
if header.endswith(b"\r\n\r\n") or not data:
break
except Exception:
pass
self._header = header
self.header_list = header.split(b"\r\n")
self._host = None
self._port = None
def get_method(self):
if self._method is None and b" " in self._header:
self._method = self._header[:self._header.index(b" ")]
return self._method
def get_host_info(self):
if self._host is None:
method = self.get_method()
line = self.header_list[0].decode("utf8", "ignore") if self.header_list else ""
if method == b"CONNECT":
host = line.split(" ")[1] if len(line.split(" ")) > 1 else ""
host, port = (host.split(":") + [443])[:2] if ":" in host else (host, 443)
else:
host = ""
for i in self.header_list:
if i.startswith(b"Host:"):
parts = i.split(b" ")
if len(parts) >= 2:
host = parts[1].decode("utf8", "ignore")
break
if not host and "/" in line:
host = line.split("/")[2]
host, port = (host.split(":") + [80])[:2] if ":" in host else (host, 80)
self._host = host
try:
self._port = int(port)
except ValueError:
self._port = 80
return self._host, self._port
@property
def data(self):
return self._header
def is_ssl(self):
return self.get_method() == b"CONNECT"
def _relay(s1, s2):
"""单向转发 s1 -> s2,直到任一端断开。"""
try:
while True:
data = s1.recv(4096)
if not data:
return
s2.sendall(data)
except Exception:
pass
finally:
try:
s1.shutdown(socket.SHUT_RD)
except Exception:
pass
def _pipe(a, b):
"""双向转发 a <-> b(两个单向线程)。"""
t1 = threading.Thread(target=_relay, args=(a, b), daemon=True)
t2 = threading.Thread(target=_relay, args=(b, a), daemon=True)
t1.start()
t2.start()
t1.join()
t2.join()
# 上游代理类型 -> PySocks 常量(延迟到 _connect_via_proxy 内 import)
_UPSTREAM_TYPES = {
"socks5": "PROXY_TYPE_SOCKS5",
"socks4": "PROXY_TYPE_SOCKS4",
"http": "PROXY_TYPE_HTTP",
}
def _connect_via_proxy(upstream_type, proxy_host, proxy_port,
dest_host, dest_port, timeout):
"""
通过指定上游代理建立到 dest 的连接(局部 socket,不污染全局)。
支持 socks5 / socks4 / http。
"""
import socks # PySocks
type_attr = _UPSTREAM_TYPES.get(upstream_type)
if type_attr is None:
raise ValueError("不支持的上游类型: {}".format(upstream_type))
ptype = getattr(socks, type_attr)
s = socks.socksocket()
s.set_proxy(ptype, proxy_host, proxy_port)
s.settimeout(timeout)
s.connect((dest_host, dest_port))
return s
def _socks5_handshake(client, first_byte):
"""
完成 SOCKS5 入站握手(无认证)。返回 (dest_host, dest_port) 或 None。
first_byte 为已读出的首字节(应为 0x05)。
"""
try:
# --- 1. 协商认证方式 ---
# 客户端: VER(1) NMETHODS(1) METHODS(NMETHODS)
nmethods = client.recv(1)
if not nmethods:
return None
client.recv(ord(nmethods)) # 丢弃 METHODS 列表
# 服务端: VER(1) METHOD(1) —— 0x00 = 无需认证
client.sendall(b"\x05\x00")
# --- 2. 读取请求 ---
# VER(1) CMD(1) RSV(1) ATYP(1) DST.ADDR(var) DST.PORT(2)
ver = client.recv(1)
if not ver or ver != b"\x05":
return None
cmd = client.recv(1)
client.recv(1) # RSV
atyp = client.recv(1)
if not atyp:
return None
atyp = ord(atyp)
if atyp == 1: # IPv4
addr = socket.inet_ntoa(client.recv(4))
elif atyp == 3: # 域名
ln = client.recv(1)
if not ln:
return None
addr = client.recv(ord(ln)).decode("utf-8", "ignore")
elif atyp == 4: # IPv6
addr = socket.inet_ntop(socket.AF_INET6, client.recv(16))
else:
client.sendall(b"\x05\x08\x00\x01\x00\x00\x00\x00\x00\x00") # 不支持 ATYP
return None
port_bytes = client.recv(2)
if len(port_bytes) < 2:
return None
port = (ord(port_bytes[0:1]) << 8) | ord(port_bytes[1:2])
# 仅支持 CONNECT (0x01)
if ord(cmd) != 0x01:
client.sendall(b"\x05\x07\x00\x01\x00\x00\x00\x00\x00\x00") # 不支持 CMD
return None
return addr, port
except Exception:
return None
def _socks5_reply_ok(client):
"""告知客户端 SOCKS5 连接已建立。"""
# VER(05) REP(00=成功) RSV(00) ATYP(01=IPv4) BND.ADDR(4) BND.PORT(2)
client.sendall(b"\x05\x00\x00\x01\x00\x00\x00\x00\x00\x00")
def _establish_upstream(pool, dest_host, dest_port, timeout, max_retries, upstream_type):
"""
尝试用池中代理建立到目标的连接,失败反馈给池子用于熔断。
成功返回 (server_socket, proxy_used),全失败返回 (None, None)。
"""
server = None
for _ in range(max_retries):
proxy = pool.pick()
if proxy is None:
log.warn("代理池为空,放弃")
return None, None
ph, pp = proxy.split(":")[0], int(proxy.split(":")[1])
try:
server = _connect_via_proxy(upstream_type, ph, pp,
dest_host, dest_port, timeout)
pool.mark_success(proxy)
log.info("via {}://{}:{}".format(upstream_type, ph, pp))
return server, proxy
except Exception as e:
pool.mark_fail(proxy)
log.warn("{} 失败: {}".format(proxy, e))
try:
if server:
server.close()