proxyfly engine v2: 3-source harvest, stratified testing, 2-strike health, transparent incl, latency cap
This commit is contained in:
29
app.py
29
app.py
@@ -99,7 +99,8 @@ def init_db():
|
|||||||
status TEXT, session_file TEXT, created REAL)""")
|
status TEXT, session_file TEXT, created REAL)""")
|
||||||
# proxy schema extensions (proxyfly integration)
|
# proxy schema extensions (proxyfly integration)
|
||||||
for col, ddl in [("clean", "INTEGER DEFAULT 0"), ("country", "TEXT"),
|
for col, ddl in [("clean", "INTEGER DEFAULT 0"), ("country", "TEXT"),
|
||||||
("egress_ip", "TEXT"), ("source", "TEXT")]:
|
("egress_ip", "TEXT"), ("source", "TEXT"),
|
||||||
|
("fail_count", "INTEGER DEFAULT 0")]:
|
||||||
try:
|
try:
|
||||||
con.execute(f"ALTER TABLE proxies ADD COLUMN {col} {ddl}")
|
con.execute(f"ALTER TABLE proxies ADD COLUMN {col} {ddl}")
|
||||||
except sqlite3.OperationalError:
|
except sqlite3.OperationalError:
|
||||||
@@ -379,8 +380,15 @@ def proxy_loop():
|
|||||||
con = db()
|
con = db()
|
||||||
for p in proxies:
|
for p in proxies:
|
||||||
health, lat, pub = results.get(p["id"], ("dead", None, None))
|
health, lat, pub = results.get(p["id"], ("dead", None, None))
|
||||||
con.execute("UPDATE proxies SET health=?, latency_ms=?, last_check=? WHERE id=?",
|
if health == "healthy":
|
||||||
(health, lat, time.time(), p["id"]))
|
con.execute("UPDATE proxies SET health='healthy', latency_ms=?, last_check=?, fail_count=0 WHERE id=?",
|
||||||
|
(lat, time.time(), p["id"]))
|
||||||
|
else:
|
||||||
|
# 2-strike rule: free proxies flake; only die after consecutive failures
|
||||||
|
fc = (p["fail_count"] or 0) + 1
|
||||||
|
con.execute("UPDATE proxies SET fail_count=?, last_check=? WHERE id=?", (fc, time.time(), p["id"]))
|
||||||
|
if fc >= 2:
|
||||||
|
con.execute("UPDATE proxies SET health='dead' WHERE id=?", (p["id"],))
|
||||||
# auto-assign: prefer verified-clean, fall back to any healthy
|
# auto-assign: prefer verified-clean, fall back to any healthy
|
||||||
if health == "healthy" and not p["assigned_mac"]:
|
if health == "healthy" and not p["assigned_mac"]:
|
||||||
devices = con.execute(
|
devices = con.execute(
|
||||||
@@ -562,7 +570,20 @@ def spawn_worker(job_id, count, prefix):
|
|||||||
{"newid": vmid, "name": name, "full": 0})
|
{"newid": vmid, "name": name, "full": 0})
|
||||||
if r.get("data") is None and r.get("error"):
|
if r.get("data") is None and r.get("error"):
|
||||||
raise Exception(f"clone failed: {r.get('error')}")
|
raise Exception(f"clone failed: {r.get('error')}")
|
||||||
pve("PUT", f"/nodes/pve/qemu/{vmid}/config", {"args": args, "memory": 3072, "cores": 2})
|
time.sleep(2)
|
||||||
|
cfg = {"args": args, "memory": 3072, "cores": 2}
|
||||||
|
applied = False
|
||||||
|
for _attempt in range(3):
|
||||||
|
pve("PUT", f"/nodes/pve/qemu/{vmid}/config", cfg)
|
||||||
|
time.sleep(2)
|
||||||
|
got = pve("GET", f"/nodes/pve/qemu/{vmid}/config")
|
||||||
|
cur = (got.get("data") or {}).get("args", "")
|
||||||
|
if f"59{display:02d}" in cur:
|
||||||
|
applied = True
|
||||||
|
break
|
||||||
|
time.sleep(3)
|
||||||
|
if not applied:
|
||||||
|
raise Exception(f"args did not apply for vmid {vmid}")
|
||||||
pve("POST", f"/nodes/pve/qemu/{vmid}/status/start", {})
|
pve("POST", f"/nodes/pve/qemu/{vmid}/status/start", {})
|
||||||
vmids.append(vmid)
|
vmids.append(vmid)
|
||||||
log("SPAWN", f"spawned {name} (vmid {vmid}, display {display})")
|
log("SPAWN", f"spawned {name} (vmid {vmid}, display {display})")
|
||||||
|
|||||||
118
proxy_harvest.py
118
proxy_harvest.py
@@ -15,13 +15,16 @@ from concurrent.futures import ThreadPoolExecutor
|
|||||||
|
|
||||||
DB_PATH = "/opt/android-fleet/fleet.db"
|
DB_PATH = "/opt/android-fleet/fleet.db"
|
||||||
SOURCES = [
|
SOURCES = [
|
||||||
"https://cdn.jsdelivr.net/gh/proxifly/free-proxy-list@main/proxies/all/data.json",
|
("proxifly", "https://cdn.jsdelivr.net/gh/proxifly/free-proxy-list@main/proxies/all/data.json"),
|
||||||
|
("geonode", "https://proxylist.geonode.com/api/proxy-list?limit=500&page=1&sort_by=lastChecked&sort_type=desc"),
|
||||||
|
("proxyscrape", "https://api.proxyscrape.com/v4/free-proxy-list/get?request=display_proxies&proxy_format=protocolipport&format=json"),
|
||||||
]
|
]
|
||||||
MAX_FETCH = 200 # candidates per cycle
|
MAX_FETCH = 200 # candidates per cycle
|
||||||
MAX_TEST = 120 # live-test at most this many per cycle
|
MAX_TEST = 120 # live-test at most this many per cycle
|
||||||
MAX_POOL = 40 # cap the proxies table (perf)
|
MAX_POOL = 60 # cap the proxies table (perf)
|
||||||
TEST_TIMEOUT = 7
|
TEST_TIMEOUT = 7
|
||||||
HARVEST_INTERVAL = 1500 # 25 min
|
HARVEST_INTERVAL = 480 # 8 min — free proxies churn fast, harvest often
|
||||||
|
MAX_LATENCY_MS = 4000 # drop anything slower than this (routeable only)
|
||||||
|
|
||||||
def log(msg):
|
def log(msg):
|
||||||
print(f"[harvest] {msg}", flush=True)
|
print(f"[harvest] {msg}", flush=True)
|
||||||
@@ -32,48 +35,87 @@ def dbcon():
|
|||||||
return con
|
return con
|
||||||
|
|
||||||
def fetch_candidates():
|
def fetch_candidates():
|
||||||
cands = []
|
cands = {"socks5": [], "socks4": [], "http": []}
|
||||||
for src in SOURCES:
|
for src_name, src in SOURCES:
|
||||||
try:
|
try:
|
||||||
with urllib.request.urlopen(src, timeout=45) as r:
|
req = urllib.request.Request(src, headers={
|
||||||
|
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 Chrome/120 Safari/537.36",
|
||||||
|
"Accept": "application/json"})
|
||||||
|
with urllib.request.urlopen(req, timeout=45) as r:
|
||||||
data = json.loads(r.read())
|
data = json.loads(r.read())
|
||||||
for p in data:
|
if src_name == "proxifly":
|
||||||
proto = p.get("protocol", "")
|
for p in data:
|
||||||
if proto not in ("socks5", "https", "http", "socks4"):
|
proto = p.get("protocol", "")
|
||||||
continue
|
if proto not in ("socks5", "https", "http", "socks4"):
|
||||||
if p.get("anonymity") == "transparent":
|
continue
|
||||||
continue # leaks real IP — useless for the fleet
|
# transparent = leaks real IP only on plain-HTTP; our traffic is
|
||||||
cands.append({
|
# HTTPS, so include them (mark them in the name later).
|
||||||
"host": p["ip"], "port": p["port"], "proto": proto,
|
cands.setdefault(proto, []).append({
|
||||||
"country": (p.get("geolocation") or {}).get("country", "?"),
|
"host": p["ip"], "port": p["port"], "proto": proto,
|
||||||
})
|
"country": (p.get("geolocation") or {}).get("country", "?"),
|
||||||
|
"anon": p.get("anonymity", ""),
|
||||||
|
})
|
||||||
|
elif src_name == "geonode":
|
||||||
|
for p in data.get("data", []):
|
||||||
|
protos = [x for x in (p.get("protocols") or []) if x in ("socks5", "socks4", "http", "https")]
|
||||||
|
if not protos:
|
||||||
|
continue
|
||||||
|
for proto in protos[:1]:
|
||||||
|
cands.setdefault(proto, []).append({
|
||||||
|
"host": p["ip"], "port": int(p["port"]), "proto": proto,
|
||||||
|
"country": p.get("country", "?"),
|
||||||
|
"anon": p.get("anonymityLevel", ""),
|
||||||
|
})
|
||||||
|
elif src_name == "proxyscrape":
|
||||||
|
for p in data.get("proxies", []):
|
||||||
|
proto = p.get("protocol", "")
|
||||||
|
if proto not in ("socks5", "socks4", "http"):
|
||||||
|
continue
|
||||||
|
cands.setdefault(proto, []).append({
|
||||||
|
"host": p["ip"], "port": int(p["port"]), "proto": proto,
|
||||||
|
"country": p.get("country", "?"),
|
||||||
|
"anon": "elite",
|
||||||
|
})
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
log(f"source error {src}: {e}")
|
log(f"source error {src_name}: {e}")
|
||||||
|
# STRATIFY + SHUFFLE: feeds are protocol-grouped; sampling the head would
|
||||||
|
# test only http proxies. Take an even mix across protocols.
|
||||||
|
import random as _random
|
||||||
|
out = []
|
||||||
|
for pool in cands.values():
|
||||||
|
_random.shuffle(pool)
|
||||||
|
while len(out) < MAX_FETCH and any(cands.values()):
|
||||||
|
for proto in ("socks5", "socks4", "http"):
|
||||||
|
if cands[proto]:
|
||||||
|
out.append(cands[proto].pop())
|
||||||
# dedupe by host:port
|
# dedupe by host:port
|
||||||
seen = set()
|
seen = set()
|
||||||
out = []
|
final = []
|
||||||
for c in cands:
|
for c in out:
|
||||||
k = f"{c['host']}:{c['port']}"
|
k = f"{c['host']}:{c['port']}"
|
||||||
if k not in seen:
|
if k not in seen:
|
||||||
seen.add(k)
|
seen.add(k)
|
||||||
out.append(c)
|
final.append(c)
|
||||||
return out[:MAX_FETCH]
|
return final[:MAX_FETCH]
|
||||||
|
|
||||||
def test_proxy(c):
|
def test_proxy(c):
|
||||||
"""Returns (c, egress_ip, latency_ms) or (c, None, None)."""
|
"""Returns (c, egress_ip, latency_ms) or (c, None, None). HTTPS first, HTTP fallback."""
|
||||||
proto_map = {"socks5": "socks5h://", "https": "https://", "http": "http://",
|
proto_map = {"socks5": "socks5h://", "https": "https://", "http": "http://",
|
||||||
"socks4": "socks4://"}
|
"socks4": "socks4://"}
|
||||||
url = proto_map[c["proto"]] + f"{c['host']}:{c['port']}"
|
url = proto_map[c["proto"]] + f"{c['host']}:{c['port']}"
|
||||||
import urllib.request as ur
|
import urllib.request as ur
|
||||||
opener = ur.build_opener(ur.ProxyHandler({"http": url, "https": url}))
|
opener = ur.build_opener(ur.ProxyHandler({"http": url, "https": url}))
|
||||||
t0 = time.time()
|
opener.addheaders = [("User-Agent", "Mozilla/5.0 (Linux; Android 9) AppleWebKit/537.36")]
|
||||||
try:
|
for target in ("https://api.ipify.org", "http://api.ipify.org"):
|
||||||
with opener.open("https://api.ipify.org", timeout=TEST_TIMEOUT) as r:
|
t0 = time.time()
|
||||||
ip = r.read().decode().strip()
|
try:
|
||||||
lat = int((time.time() - t0) * 1000)
|
with opener.open(target, timeout=TEST_TIMEOUT) as r:
|
||||||
return (c, ip, lat)
|
ip = r.read().decode().strip()
|
||||||
except Exception:
|
lat = int((time.time() - t0) * 1000)
|
||||||
return (c, None, None)
|
return (c, ip, lat)
|
||||||
|
except Exception:
|
||||||
|
continue
|
||||||
|
return (c, None, None)
|
||||||
|
|
||||||
def check_clean(ips):
|
def check_clean(ips):
|
||||||
"""ip-api batch: flags for each ip. Returns {ip: dict}."""
|
"""ip-api batch: flags for each ip. Returns {ip: dict}."""
|
||||||
@@ -109,6 +151,8 @@ def dnsbl_listed(ip):
|
|||||||
|
|
||||||
def prune_pool():
|
def prune_pool():
|
||||||
con = dbcon()
|
con = dbcon()
|
||||||
|
# drop dead proxifly entries every cycle (they die fast — keep the list honest)
|
||||||
|
con.execute("DELETE FROM proxies WHERE source='proxifly' AND health='dead' AND assigned_mac IS NULL")
|
||||||
rows = con.execute("SELECT id, egress_ip, source FROM proxies ORDER BY id").fetchall()
|
rows = con.execute("SELECT id, egress_ip, source FROM proxies ORDER BY id").fetchall()
|
||||||
# remove duplicates by egress ip (keep lowest id)
|
# remove duplicates by egress ip (keep lowest id)
|
||||||
seen = set()
|
seen = set()
|
||||||
@@ -151,12 +195,11 @@ def harvest():
|
|||||||
dns = dns_flags.get(c["egress_ip"], False)
|
dns = dns_flags.get(c["egress_ip"], False)
|
||||||
clean = 1 if (not f.get("hosting") and not f.get("proxy") and not dns) else 0
|
clean = 1 if (not f.get("hosting") and not f.get("proxy") and not dns) else 0
|
||||||
mobile = 1 if f.get("mobile") else 0
|
mobile = 1 if f.get("mobile") else 0
|
||||||
# keep only clean ones + a few fallbacks, skip dirty datacenter spam
|
if c["latency_ms"] and c["latency_ms"] > MAX_LATENCY_MS:
|
||||||
if clean == 0:
|
continue # too slow to route through
|
||||||
why = "hosting" if f.get("hosting") else ("proxy" if f.get("proxy") else ("dnsbl" if dns else "unknown"))
|
# BANK EVERYTHING ALIVE: grade it, route clean-first. The list stays populated.
|
||||||
log(f"rejected {c['egress_ip']} ({why})")
|
suffix = "-T" if c.get("anon") == "transparent" else ""
|
||||||
continue
|
name = f"proxifly-{c['country']}-{c['egress_ip'].split('.')[-2]}{suffix}"
|
||||||
name = f"proxifly-{c['country']}-{c['egress_ip'].split('.')[-2]}"
|
|
||||||
con.execute("""INSERT INTO proxies(name, host, port, proto, health, latency_ms,
|
con.execute("""INSERT INTO proxies(name, host, port, proto, health, latency_ms,
|
||||||
assigned_mac, last_check, clean, country, egress_ip, source)
|
assigned_mac, last_check, clean, country, egress_ip, source)
|
||||||
VALUES(?,?,?,?,'healthy',?,NULL,?,?,?,?,?)
|
VALUES(?,?,?,?,'healthy',?,NULL,?,?,?,?,?)
|
||||||
@@ -182,13 +225,14 @@ def assign_best():
|
|||||||
have = con.execute("SELECT * FROM proxies WHERE assigned_mac=?", (d["mac"],)).fetchone()
|
have = con.execute("SELECT * FROM proxies WHERE assigned_mac=?", (d["mac"],)).fetchone()
|
||||||
if have:
|
if have:
|
||||||
continue # already assigned — keep stable
|
continue # already assigned — keep stable
|
||||||
|
# assignment preference: verified-clean first, then any live proxy
|
||||||
best = con.execute("""SELECT * FROM proxies
|
best = con.execute("""SELECT * FROM proxies
|
||||||
WHERE clean=1 AND health='healthy' AND assigned_mac IS NULL
|
WHERE clean=1 AND health='healthy' AND assigned_mac IS NULL
|
||||||
ORDER BY latency_ms ASC LIMIT 1""").fetchone()
|
ORDER BY latency_ms ASC LIMIT 1""").fetchone()
|
||||||
if not best:
|
if not best:
|
||||||
best = con.execute("""SELECT * FROM proxies
|
best = con.execute("""SELECT * FROM proxies
|
||||||
WHERE health='healthy' AND assigned_mac IS NULL
|
WHERE health='healthy' AND assigned_mac IS NULL
|
||||||
ORDER BY latency_ms ASC LIMIT 1""").fetchone()
|
ORDER BY clean DESC, latency_ms ASC LIMIT 1""").fetchone()
|
||||||
if not best:
|
if not best:
|
||||||
continue
|
continue
|
||||||
con.execute("UPDATE proxies SET assigned_mac=? WHERE id=?", (d["mac"], best["id"]))
|
con.execute("UPDATE proxies SET assigned_mac=? WHERE id=?", (d["mac"], best["id"]))
|
||||||
|
|||||||
Reference in New Issue
Block a user