"""图形选股本地代理:解决浏览器直连东财 CORS 失败问题。零依赖,直接运行。 运行: python server.py -> 浏览器打开 http://127.0.0.1:8000/index.html """ import json, urllib.request, urllib.parse, os CURL = "curl.exe" if os.name == "nt" else "curl" from http.server import SimpleHTTPRequestHandler, HTTPServer EA_HIS = "https://push2his.eastmoney.com/api/qt/stock/kline/get" EA_LIST = "https://push2.eastmoney.com/api/qt/clist/get" TX = "https://qt.gtimg.cn/q=" EA_DEAD = {"his": False, "list": False} import time as _time HEART = {"beat": _time.time(), "seen": False} def _touch(): HEART["beat"] = _time.time() HEART["seen"] = True def get(url): req = urllib.request.Request(url, headers={"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120 Safari/537.36", "Referer": "https://quote.eastmoney.com/", "Accept": "*/*", "Connection": "keep-alive"}) with urllib.request.urlopen(req, timeout=5) as r: return r.read() def tx_kline(code, beg, end): """优先腾讯前复权(qfq),失败回退搜狐不复权;无数据返回空klines而非抛错""" import subprocess, datetime today = datetime.date.today().strftime("%Y%m%d") if end > today: end = today if beg > end: return json.dumps({"data": {"name": code, "klines": []}}).encode() pre = "sh" if code.startswith("6") else "sz" try: import datetime as _dt d1 = _dt.datetime.strptime(beg, "%Y%m%d"); d2 = _dt.datetime.strptime(end, "%Y%m%d") cnt = max(30, min(1000, (d2 - d1).days + 15)) url = "https://ifzq.gtimg.cn/appstock/app/fqkline/get?param=%s%s,day,,,%d,qfq" % (pre, code, cnt) raw = subprocess.check_output([CURL, "-s", "-m", "20", "-A", "Mozilla/5.0", url], timeout=25).decode("utf-8", "ignore") j = json.loads(raw) rows = j["data"][pre + code]["qfqday"] kl = [] for d in rows: date = d[0].replace("-", "") if date < beg or date > end: continue o, c, h, l, v = d[1], d[2], d[3], d[4], d[5] kl.append("%s,%s,%s,%s,%s,%s,0" % (date, o, c, h, l, v)) if kl: return json.dumps({"data": {"name": code, "klines": kl}}).encode() except Exception: pass pre2 = "cn_%s" % code url2 = "https://q.stock.sohu.com/hisHq?code=%s&start=%s&end=%s" % (pre2, beg, end) try: raw = subprocess.check_output([CURL, "-s", "-m", "20", "-A", "Mozilla/5.0", url2], timeout=25).decode("utf-8", "ignore") j = json.loads(raw) rows = j[0]["hq"] if j and isinstance(j, list) and j[0].get("hq") else [] except Exception: return json.dumps({"data": {"name": code, "klines": []}}).encode() kl = [] name = code for d in rows: date = d[0].replace("-", "") if date < beg or date > end: continue o, c, l, h = d[1], d[2], d[5], d[6] kl.append("%s,%s,%s,%s,%s,%s,0" % (date, o, c, h, l, d[8].replace(",", ""))) kl.reverse() return json.dumps({"data": {"name": name, "klines": kl}}).encode() import subprocess def curl(url, ua="Mozilla/5.0"): return subprocess.check_output([CURL, "-s", "-m", "20", "-A", ua, url], timeout=25) # ---- 登录:硬件ID / 验证码 / MQTT(EMQX) 校验 ---- import socket, struct, random, time, threading CAPTCHA = {} CAPTCHA_LOCK = threading.Lock() CAPTCHA_CHARS = "23456789" MQTT_HOST = os.environ.get("MQTT_HOST", "8.137.151.237") MQTT_PORT = int(os.environ.get("MQTT_PORT", "1883")) def machine_id(): if os.name == "nt": try: import re out = subprocess.check_output(["reg", "query", r"HKLM\SOFTWARE\Microsoft\Cryptography", "/v", "MachineGuid"], timeout=5).decode("gbk", "ignore") m = re.search(r"MachineGuid\s+REG_SZ\s+([0-9a-fA-F\-]+)", out) if m: return m.group(1) except Exception: pass try: import uuid return "%012x" % uuid.getnode() except Exception: return "unknown" def _captcha_svg(code): W, H, gap, pad = 22, 36, 12, 14 n = len(code) width = pad * 2 + n * W + (n - 1) * gap height = pad * 2 + H P = ['' % (width, height, width, height)] P.append('') SEG = {'a': (0, 0, W, 0), 'g': (0, H / 2, W, H / 2), 'd': (0, H, W, H), 'f': (0, 0, 0, H / 2), 'b': (W, 0, W, H / 2), 'e': (0, H / 2, 0, H), 'c': (W, H / 2, W, H)} DMAP = {'2': 'abged', '3': 'abgcd', '4': 'fgbc', '5': 'afgcd', '6': 'afgecd', '7': 'abc', '8': 'abcdefg', '9': 'abcdfg'} colors = ['#e5484d', '#4da3ff', '#f5a623', '#26a69a', '#b266ff'] for i, ch in enumerate(code): x0 = pad + i * (W + gap) col = random.choice(colors) ang = random.randint(-14, 14) P.append('' % (x0, pad + random.randint(-3, 3), ang, W / 2, H / 2, col)) for s in DMAP.get(ch, ''): x1, y1, x2, y2 = SEG[s] P.append('' % (x1, y1, x2, y2)) P.append('') for _ in range(5): P.append('' % (random.randint(0, width), random.randint(0, height), random.randint(0, width), random.randint(0, height))) for _ in range(18): P.append('' % (random.randint(0, width), random.randint(0, height))) P.append('') return "".join(P) def new_captcha(): code = "".join(random.choice(CAPTCHA_CHARS) for _ in range(4)) token = "%016x" % random.getrandbits(64) with CAPTCHA_LOCK: now = time.time() for k in list(CAPTCHA.keys()): if now - CAPTCHA[k][1] > 300: del CAPTCHA[k] CAPTCHA[token] = (code, now) return token, _captcha_svg(code) def check_captcha(token, ans): with CAPTCHA_LOCK: rec = CAPTCHA.get(token) if not rec: return False code, ts = rec if time.time() - ts > 300: CAPTCHA.pop(token, None) return False if str(ans).strip().upper() == code.upper(): CAPTCHA.pop(token, None) return True return False def mqtt_auth(user, pw, client_id, host=None, port=None, timeout=8): if host is None: host = MQTT_HOST if port is None: port = MQTT_PORT try: def enc(s): b = s.encode("utf-8") return struct.pack("!H", len(b)) + b flags = 0x02 if user: flags |= 0x80 if pw: flags |= 0x40 payload = enc(client_id) + enc(user) + enc(pw) var = enc("MQTT") + bytes([0x04, flags]) + struct.pack("!H", 30) body = var + payload rl = len(body) rem = b"" while True: d = rl % 128 rl //= 128 if rl > 0: rem += bytes([d | 0x80]) else: rem += bytes([d]) break pkt = b"\x10" + rem + body s = socket.create_connection((host, port), timeout=timeout) try: s.sendall(pkt) resp = s.recv(4) finally: s.close() if len(resp) >= 4 and resp[0] == 0x20: return resp[3] if len(resp) >= 2 and (resp[0] & 0xF0) == 0x20: return resp[1] if len(resp) == 2 else resp[3] return -2 except Exception: return -1 # ---- 远程升级:经 MQTT(EMQX) 拉取版本公告 + 分片文件 ---- UP_PREFIX = os.environ.get("MQTT_UP_PREFIX", "gp") UP_KEY = os.environ.get("UPGRADE_KEY", "qqq7079291") CHUNK_RAW = 24 * 1024 UP_STAT = {"pct": 0, "msg": ""} UP_DIR = os.environ.get("UPGRADE_DIR", "up") PUBLIC_BASE = os.environ.get("PUBLIC_BASE", "https://up.youxue.space").rstrip("/") HUB_KEY = os.environ.get("HUB_KEY", "") HUB_MQTT_USER = os.environ.get("HUB_MQTT_USER", "") HUB_MQTT_PASS = os.environ.get("HUB_MQTT_PASS", "") def _ustat(pct, msg): try: UP_STAT["pct"] = pct UP_STAT["msg"] = msg except Exception: pass def _mstr(s): b = s.encode("utf-8") return struct.pack("!H", len(b)) + b def _mhead(typ, flags, body): rl = len(body) rem = b"" while True: d = rl % 128 rl //= 128 if rl > 0: rem += bytes([d | 0x80]) else: rem += bytes([d]) break return bytes([(typ << 4) | flags]) + rem + body def mqtt_session(user, pw, cid, timeout=10): s = socket.create_connection((MQTT_HOST, MQTT_PORT), timeout=timeout) s.settimeout(timeout) try: flags = 0x02 if user: flags |= 0x80 if pw: flags |= 0x40 body = _mstr("MQTT") + bytes([0x04, flags]) + struct.pack("!H", 60) + _mstr(cid) + _mstr(user) + _mstr(pw) s.sendall(_mhead(1, 0, body)) typ, _flags, payload = _mread(s, timeout) if typ != 2 or len(payload) < 2: raise RuntimeError("connack-missing") if payload[1] != 0: raise RuntimeError("mqtt-rc-%d" % payload[1]) return s except Exception: try: s.close() except Exception: pass raise def _mread(s, timeout): import time as _t buf = getattr(s, "_mbuf", b"") end = _t.time() + timeout def _need(n): nonlocal buf while len(buf) < n: left = end - _t.time() if left <= 0: raise TimeoutError("mqtt-timeout") s.settimeout(left) r = s.recv(65536) if not r: raise ConnectionError("mqtt-closed") buf += r _need(2) typ = buf[0] >> 4 flags = buf[0] & 0x0F mult = 1 rl = 0 i = 1 while True: if i >= len(buf): _need(i + 1) d = buf[i] i += 1 rl += (d & 127) * mult mult *= 128 if not (d & 128): break _need(i + rl) body = buf[i:i + rl] try: s._mbuf = buf[i + rl:] except Exception: pass return typ, flags, body def mqtt_sub(s, topic, qos=0): import random as _r pid = _r.randint(1, 60000) body = struct.pack("!H", pid) + _mstr(topic) + bytes([qos]) s.sendall(_mhead(8, 2, body)) def mqtt_pub(s, topic, payload, retain=False, qos=0): if isinstance(payload, str): payload = payload.encode("utf-8") body = _mstr(topic) + payload s.sendall(_mhead(3, ((qos & 3) << 1) | (1 if retain else 0), body)) def _mping(s): try: s.sendall(_mhead(12, 0, b"")) except Exception: pass def mqtt_collect(s, topics, timeout=10, stop=None, prog=None, qos=1): import time as _t for t in topics: mqtt_sub(s, t, qos=qos) out = {} denied = [] pend = list(topics) end = _t.time() + timeout last_ping = _t.time() while True: if stop: try: if stop(out): break except Exception: pass left = end - _t.time() if left <= 0: break if _t.time() - last_ping > 20: _mping(s) last_ping = _t.time() try: typ, flags, body = _mread(s, min(left, 25)) except Exception: break if typ == 3 and len(body) > 2: tl = struct.unpack("!H", body[:2])[0] pq = (flags >> 1) & 3 off = 2 + tl if pq > 0 and len(body) >= off + 2: try: s.sendall(_mhead(4, 0, body[off:off + 2])) except Exception: pass off += 2 topic = body[2:2 + tl].decode("utf-8", "ignore") payload = body[off:] out.setdefault(topic, []).append(payload) if prog: try: prog(out) except Exception: pass elif typ == 9 and len(body) > 2: for rc in body[2:]: t = pend.pop(0) if pend else "" if rc == 0x80 and t and t not in denied: denied.append(t) elif typ == 13: _mping(s) last_ping = _t.time() out["_denied"] = denied return out def parse_ver(v): import re as _re v = (v or "").strip() if v[:1].lower() == "v": v = v[1:] m = _re.match(r"(\d+)\.(\d+)\.(\d+)([a-z]*)", v, _re.I) if not m: return (0, 0, 0, v) return (int(m.group(1)), int(m.group(2)), int(m.group(3)), (m.group(4) or "").lower()) def local_ver(): try: import re as _re h = open("index.html", encoding="utf-8").read() m = _re.search(r'id="ver">([^<]+)<', h) if m: return m.group(1).strip() except Exception: pass return "v0" def up_manifest_topic(): return UP_PREFIX + "/upgrade/manifest" def up_sig_payload(m): parts = [str(m.get("latest", "")), str(m.get("notes", ""))] for f in m.get("files", []) or []: parts.append("%s|%s|%s|%s" % (f.get("name", ""), f.get("size", ""), f.get("sha256", ""), f.get("msgs", ""))) return "|".join(parts) def upgrade_check(user, pw, timeout=10): import json as _j, hashlib as _h cid = "gp-up-%016x" % random.getrandbits(64) s = mqtt_session(user, pw, cid, timeout=timeout) try: mt = up_manifest_topic() got = mqtt_collect(s, [mt], timeout=timeout, stop=lambda o: bool(o.get(mt))) finally: try: s.close() except Exception: pass msgs = got.get(up_manifest_topic()) or [] if not msgs: return {"ok": True, "has": False, "current": local_ver()} try: m = _j.loads(msgs[0].decode("utf-8", "ignore")) except Exception: return {"ok": False, "msg": "版本公告解析失败"} cur = local_ver() latest = str(m.get("latest", "")) need = parse_ver(latest) > parse_ver(cur) verified = True if UP_KEY: sig = str(m.get("sig", "")) expect = _h.sha256((UP_KEY + "|" + up_sig_payload(m)).encode("utf-8")).hexdigest() verified = (sig == expect) if need and not verified: return {"ok": False, "msg": "版本签名校验失败,拒绝升级", "current": cur, "latest": latest} return {"ok": True, "has": True, "current": cur, "latest": latest, "need": need, "notes": m.get("notes", ""), "verified": verified, "manifest": m, "files": [{"name": f.get("name"), "size": f.get("size")} for f in (m.get("files") or [])]} def upgrade_apply(user, pw, timeout=120): import json as _j, hashlib as _h, base64 as _b, time as _t, shutil as _sh _ustat(5, "检查版本公告…") chk = None for _att in range(3): try: chk = upgrade_check(user, pw, timeout=10) except Exception as e: chk = {"ok": False, "msg": "连接升级服务失败:%s" % str(e)} if chk.get("ok") and chk.get("has"): break if _att < 2: try: _t.sleep(2) except Exception: pass if chk is None: return {"ok": False, "msg": "版本公告获取失败,请重试"} if not chk.get("ok"): return chk if not chk.get("ok"): return chk if not chk.get("has"): return {"ok": False, "msg": "暂无版本公告"} if not chk.get("need"): return {"ok": True, "applied": False, "msg": "已是最新版本"} latest = chk["latest"] m = chk.get("manifest") or {} files = m.get("files") or [] if not files: return {"ok": False, "msg": "升级包为空"} _ustat(20, "版本公告就绪,开始下载…") base = "%s/up/%s/" % (UP_PREFIX, latest) want = {} for f in files: want[base + f["name"] + "/"] = int(f.get("msgs", 0)) total = sum(want.values()) or 1 merged = {} all_denied = [] def _merge(o): for t, arr in o.items(): if t == "_denied": for d in arr: if d not in all_denied: all_denied.append(d) continue merged.setdefault(t, []).extend(arr) def _have(pre): c = 0 for t, arr in merged.items(): if t.startswith(pre): c += len(arr) return c def _all_in(o): _merge(o) for pre, n in want.items(): if _have(pre) < n: return False return True def _prog(o): _merge(o) c = sum(min(_have(pre), n) for pre, n in want.items()) _ustat(25 + int(55 * min(c, total) / total), "下载升级包 %d/%d…" % (min(c, total), total)) topics = ["%s/up/%s/%s/#" % (UP_PREFIX, latest, f["name"]) for f in files] s2 = mqtt_session(user, pw, "gp-dl-%016x" % random.getrandbits(64), timeout=15) try: mqtt_collect(s2, topics, timeout=25, stop=_all_in, prog=_prog) for rnd in range(2): miss = [] for f in files: pre = base + f["name"] + "/" n = int(f.get("msgs", 0)) got_i = set() for t, arr in merged.items(): if t.startswith(pre): for p in arr: try: got_i.add(int(_j.loads(p.decode("utf-8", "ignore"))["i"])) except Exception: pass for i in range(n): if i not in got_i: miss.append("%s%04d" % (pre, i)) if not miss: break _ustat(80, "补齐缺失分片(第%d次)…" % (rnd + 1)) mqtt_collect(s2, miss, timeout=12, prog=_prog) finally: try: s2.close() except Exception: pass if all_denied: _ustat(0, "订阅被拒绝") return {"ok": False, "msg": "账号无权订阅升级主题,请联系管理员开权限(%s)" % all_denied[0]} for f in files: name = f["name"] blobs = {} for f in files: name = f["name"] if f.get("url"): _ustat(30, "HTTP下载 %s…" % name) tmp = name + ".dl" try: http_dl(f["url"], tmp) with open(tmp, "rb") as fp: raw = fp.read() finally: try: os.remove(tmp) except Exception: pass if _h.sha256(raw).hexdigest() != f.get("sha256"): _ustat(0, "校验失败") return {"ok": False, "msg": "文件 %s 校验失败" % name} blobs[name] = raw continue pre = base + name + "/" parts = {} for t, arr in merged.items(): if t.startswith(pre): for p in arr: try: c = _j.loads(p.decode("utf-8", "ignore")) parts[int(c["i"])] = c["data"] except Exception: pass n = int(f.get("msgs", 0)) if len(parts) < n: return {"ok": False, "msg": "文件 %s 分片不全(%d/%d),稍后重试" % (name, len(parts), n)} raw = b"".join(_b.b64decode(parts[i]) for i in range(n)) if _h.sha256(raw).hexdigest() != f.get("sha256"): _ustat(0, "校验失败") return {"ok": False, "msg": "文件 %s 校验失败" % name} blobs[name] = raw _ustat(88, "校验通过,备份并安装…") ts = _t.strftime("%Y%m%d-%H%M%S") bdir = os.path.join("backup", ts) os.makedirs(bdir, exist_ok=True) applied = [] for name, raw in blobs.items(): if name.lower().endswith(".zip"): import zipfile as _z, io as _io try: zf = _z.ZipFile(_io.BytesIO(raw)) except Exception: return {"ok": False, "msg": "压缩包 %s 解析失败" % name} members = [m for m in zf.namelist() if m and not m.endswith("/")] safe = [] for m in members: p = os.path.normpath(m).replace("\\", "/") if p.startswith("/") or p.startswith("..") or "/../" in p: return {"ok": False, "msg": "压缩包内含非法路径:%s" % m} safe.append(p) with open(os.path.join(bdir, os.path.basename(name)), "wb") as fp: fp.write(raw) for orig, p in zip(members, safe): if os.path.exists(p): dst = os.path.join(bdir, p + ".bak") try: dd = os.path.dirname(dst) if dd: os.makedirs(dd, exist_ok=True) except Exception: pass _sh.copy2(p, dst) try: d = os.path.dirname(p) if d: os.makedirs(d, exist_ok=True) with open(p, "wb") as fp: fp.write(zf.read(orig)) except Exception as e: return {"ok": False, "msg": "解压 %s 失败:%s" % (p, e)} applied.extend(safe) else: if os.path.exists(name): _sh.copy2(name, os.path.join(bdir, os.path.basename(name) + ".bak")) with open(name, "wb") as fp: fp.write(raw) applied.append(name) _ustat(100, "安装完成") return {"ok": True, "applied": True, "latest": latest, "files": sorted(applied), "backup": bdir} def _self_restart(): import sys as _sys, subprocess as _sp try: exe = _sys.executable or "py" _sp.Popen([exe, os.path.abspath("server.py")], cwd=os.path.dirname(os.path.abspath("server.py")) or ".", creationflags=0x8 if os.name == "nt" else 0, close_fds=True) except Exception: pass os._exit(0) # ---- 中转发布(hub):收作者推送 -> 落盘 -> MQTT发公告+分片 ---- def hub_publish_package(version, notes, blobs): import json as _j, hashlib as _h, base64 as _b, re as _re if not _re.match(r"^v\d+\.\d+\.\d+[a-z]*$", version, _re.I): raise ValueError("版本号格式错误") if not blobs: raise ValueError("空文件包") vdir = os.path.join(UP_DIR, version) os.makedirs(vdir, exist_ok=True) entries = [] for name, raw in blobs.items(): if not name or name.startswith("/") or ".." in name.replace("\\", "/") or len(raw) > 50 * 1024 * 1024: raise ValueError("非法文件名或超大:%s" % name) with open(os.path.join(vdir, name), "wb") as fp: fp.write(raw) entries.append({"name": name, "size": len(raw), "sha256": _h.sha256(raw).hexdigest()}) s = mqtt_session(HUB_MQTT_USER, HUB_MQTT_PASS, "gp-hub-%016x" % random.getrandbits(64), timeout=15) try: for e in entries: raw = blobs[e["name"]] chunks = [raw[i:i + CHUNK_RAW] for i in range(0, len(raw), CHUNK_RAW)] or [b""] e["msgs"] = len(chunks) e["url"] = "%s/up/%s/%s" % (PUBLIC_BASE, version, e["name"]) base = "%s/up/%s/%s/" % (UP_PREFIX, version, e["name"]) for i, c in enumerate(chunks): mqtt_pub(s, "%s%04d" % (base, i), _j.dumps({"i": i, "n": len(chunks), "data": _b.b64encode(c).decode()}), retain=True, qos=1) manifest = {"latest": version, "notes": notes, "files": entries} if UP_KEY: import hashlib as _h2 manifest["sig"] = _h2.sha256((UP_KEY + "|" + up_sig_payload(manifest)).encode("utf-8")).hexdigest() mqtt_pub(s, up_manifest_topic(), _j.dumps(manifest, ensure_ascii=False), retain=True, qos=1) finally: try: s.close() except Exception: pass return manifest # ---- HTTP下载(断点续传) ---- def http_dl(url, dst, p0=25, p1=80): import urllib.request as _u, urllib.error as _ue have = os.path.getsize(dst) if os.path.exists(dst) else 0 total = 0 for _ in range(3): req = _u.Request(url, headers={"User-Agent": "GP-Updater"}) if have > 0: req.add_header("Range", "bytes=%d-" % have) try: r = _u.urlopen(req, timeout=30) except _ue.HTTPError as e: if e.code == 416 and have > 0: return have have = 0 continue code = getattr(r, "status", 200) if code == 206: cr = r.headers.get("Content-Range", "") try: total = int(cr.split("/")[-1]) except Exception: total = 0 mode = "ab" else: try: total = int(r.headers.get("Content-Length") or 0) except Exception: total = 0 have = 0 mode = "wb" with open(dst, mode) as fp: while True: buf = r.read(256 * 1024) if not buf: break fp.write(buf) have += len(buf) if total: _ustat(p0 + int((p1 - p0) * min(have, total) / total), "下载升级包 %d/%dKB…" % (have // 1024, total // 1024)) try: r.close() except Exception: pass if total and have >= total: break if not total: break return have class H(SimpleHTTPRequestHandler): def do_GET(self): p = urllib.parse.urlparse(self.path) q = urllib.parse.parse_qs(p.query) one = lambda k, d="": q.get(k, [d])[0] _touch() try: if p.path == "/api/ping": data = json.dumps({"ok": True}).encode() self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/upgrade/status": data = json.dumps({"ok": True, "pct": UP_STAT.get("pct", 0), "msg": UP_STAT.get("msg", "")}, ensure_ascii=False).encode("utf-8") self.send_response(200); self.send_header("Content-Type", "application/json; charset=utf-8"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/hwid": data = json.dumps({"hwid": machine_id()}).encode() self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/captcha": token, svg = new_captcha() data = json.dumps({"token": token, "svg": svg}).encode() self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/kline": code, beg, end = one("code"), one("beg"), one("end") data = None try: data = tx_kline(code, beg, end) j = json.loads(data.decode()) if not j.get("data", {}).get("klines"): data = None except Exception: data = None if data is None and not EA_DEAD["his"]: try: secid = ("1." if code.startswith("6") else "0.") + code url = EA_HIS + "?" + urllib.parse.urlencode({"secid": secid, "klt": 101, "fqt": 1, "lmt": 1000, "beg": beg, "end": end, "fields1": "f1,f2,f3,f4,f5", "fields2": "f51,f52,f53,f54,f55,f56,f57"}) data = get(url) if b"klines" not in data: raise RuntimeError("eastmoney empty") except Exception: EA_DEAD["his"] = True data = None if data is None: data = json.dumps({"data": {"name": code, "klines": []}}).encode() self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/bkl": import subprocess, concurrent.futures as _cf codes = one("codes").split(",") beg, end = one("beg"), one("end") def one_k(code): code = code.strip() if not code: return (code, None) try: raw = tx_kline(code, beg, end) import json as _j j = _j.loads(raw.decode()) kl = j["data"]["klines"] if not kl: return (code, None) return (code, {"name": j["data"]["name"], "klines": kl}) except Exception: return (code, None) out = {} with _cf.ThreadPoolExecutor(max_workers=80) as ex: for code, v in ex.map(one_k, codes): if v: out[code] = v # 腾讯源并发时偶发丢码,对失败的再重试一次(低并发) miss = [c.strip() for c in codes if c.strip() and c.strip() not in out] for _round in range(2): if not miss: break still = [] for c in miss: try: raw = tx_kline(c, beg, end) j = json.loads(raw.decode()) kl = j["data"]["klines"] if kl: out[c] = {"name": j["data"]["name"], "klines": kl} else: still.append(c) except Exception: still.append(c) import time as _t _t.sleep(0.05) miss = still import json as _jj body = _jj.dumps(out).encode() self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(body); return if p.path == "/api/clist": url = EA_LIST + "?" + urllib.parse.urlencode({"pn": one("pn","1"), "pz": one("pz","500"), "po": 1, "np": 1, "fltt": 2, "invt": 2, "fid": "f3", "fs": one("fs"), "fields": "f12,f14"}) try: data = get(url) if b"diff" not in data: raise RuntimeError("empty") except Exception: data = curl(url) self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/idx": code = one("code", "000001") import datetime as _dt, threading as _th global _idx_cache, _idx_lock try: _idx_cache except NameError: _idx_cache = {} _idx_lock = _th.Lock() today = _dt.date.today().strftime("%Y-%m-%d") hit = _idx_cache.get(code) if hit and hit.get("day") == today and hit.get("klines"): data = json.dumps({"data": {"name": code, "klines": hit["klines"]}}).encode() self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return with _idx_lock: hit = _idx_cache.get(code) if hit and hit.get("day") == today and hit.get("klines"): data = json.dumps({"data": {"name": code, "klines": hit["klines"]}}).encode() else: zmap = {"000001": "zs_000001", "399001": "zs_399001", "399006": "zs_399006", "000688": "zs_000688"} zs = zmap.get(code, "zs_000001") end = _dt.date.today().strftime("%Y%m%d") beg = (_dt.date.today() - _dt.timedelta(days=400)).strftime("%Y%m%d") url = "https://q.stock.sohu.com/hisHq?code=%s&start=%s&end=%s" % (zs, beg, end) try: raw = curl(url).decode("utf-8", "ignore") j = json.loads(raw) rows = j[0]["hq"] if j and isinstance(j, list) and j[0].get("hq") else [] kl = [] for d in rows: date = d[0].replace("-", "") o, c, l, h = d[1], d[2], d[5], d[6] kl.append("%s,%s,%s,%s,%s,%s,%s" % (date, o, c, h, l, d[7] if len(d) > 7 else 0, d[8] if len(d) > 8 else 0)) kl.reverse() if not kl and hit and hit.get("klines"): kl = hit["klines"] _idx_cache[code] = {"day": today, "klines": kl} data = json.dumps({"data": {"name": code, "klines": kl}}).encode() except Exception: kl = (hit.get("klines") if hit else []) or [] data = json.dumps({"data": {"name": code, "klines": kl}}).encode() self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/lhb": base = "https://datacenter.eastmoney.com/api/data/v1/get" params = {"sortColumns": "TRADE_DATE,SECURITY_CODE", "sortTypes": "-1,-1", "pageSize": one("pz", "30"), "pageNumber": one("pn", "1"), "reportName": "RPT_DAILYBILLBOARD_DETAILS", "columns": "ALL", "source": "WEB", "client": "WEB"} url = base + "?" + urllib.parse.urlencode(params) try: data = get(url) except Exception: data = curl(url) self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/breadth": import datetime as _dt2 global _breadth_cache try: _breadth_cache except NameError: _breadth_cache = {"day": "", "data": None} today = _dt2.date.today().strftime("%Y-%m-%d") if _breadth_cache.get("day") == today and _breadth_cache.get("data"): data = json.dumps(_breadth_cache["data"]).encode() self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return import concurrent.futures as _cf2 codes = [] for a, b, pre in [(600000, 601999, "sh"), (603000, 605999, "sh"), (688000, 688999, "sh"), (1, 3999, "sz"), (300000, 301999, "sz")]: for i in range(a, b + 1): c = str(i).zfill(6) if pre == "sz" and a == 1 and i > 3999: break codes.append(pre + c) codes = [c for c in codes if not (c[2:] .startswith("6889") and False)] def one_batch(lst): try: u = "https://hq.sinajs.cn/list=" + ",".join(lst) req = urllib.request.Request(u, headers={"User-Agent": "Mozilla/5.0", "Referer": "https://finance.sina.com.cn/"}) raw = urllib.request.urlopen(req, timeout=20).read().decode("gbk", "ignore") return raw except Exception: return "" up = dn = fl = zt = dt_ = z20 = d20 = total = 0 amount = 0.0 buckets = {"跌停": 0, "跌9~7": 0, "跌7~5": 0, "跌5~3": 0, "跌3~1": 0, "跌1~0": 0, "平": 0, "涨0~1": 0, "涨1~3": 0, "涨3~5": 0, "涨5~7": 0, "涨7~9": 0, "涨停": 0} batches = [codes[i:i + 400] for i in range(0, len(codes), 400)] with _cf2.ThreadPoolExecutor(max_workers=10) as ex: raws = list(ex.map(one_batch, batches)) for raw in raws: for line in raw.split("\n"): q = line.find('="') if q < 0: continue sym = line[8:q] payload = line[q + 2:].strip().strip('";') if not payload: continue p = payload.split(",") if len(p) < 4: continue try: prev = float(p[2] or 0) cur = float(p[3] or 0) except Exception: continue if prev <= 0 or cur <= 0: continue total += 1 cp = (cur - prev) / prev * 100 try: amount += float(p[9] or 0) except Exception: pass code6 = sym[2:] if len(sym) == 8 else sym is20 = code6.startswith("300") or code6.startswith("301") or code6.startswith("688") lim = 19.5 if is20 else 9.5 if cp >= lim - 0.05: zt += 1 if is20: z20 += 1 buckets["涨停"] += 1 elif cp <= -lim + 0.05: dt_ += 1 if is20: d20 += 1 buckets["跌停"] += 1 elif cp >= 7: buckets["涨7~9"] += 1 elif cp >= 5: buckets["涨5~7"] += 1 elif cp >= 3: buckets["涨3~5"] += 1 elif cp >= 1: buckets["涨1~3"] += 1 elif cp > 0.001: buckets["涨0~1"] += 1 elif cp <= -7: buckets["跌9~7"] += 1 elif cp <= -5: buckets["跌7~5"] += 1 elif cp <= -3: buckets["跌5~3"] += 1 elif cp <= -1: buckets["跌3~1"] += 1 elif cp < -0.001: buckets["跌1~0"] += 1 else: buckets["平"] += 1 if cp > 0.001: up += 1 elif cp < -0.001: dn += 1 else: fl += 1 out = {"day": today, "total": total, "up": up, "down": dn, "flat": fl, "zt": zt, "dt": dt_, "zt20": z20, "dt20": d20, "amount": amount, "buckets": buckets} _breadth_cache = {"day": today, "data": out} self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(json.dumps(out).encode()); return if p.path == "/api/hot": self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(json.dumps({"data": [], "note": "hot-by-local"}).encode()); return if p.path == "/api/theme": codes = [c.strip() for c in one("codes").split(",") if c.strip()] import concurrent.futures as _cf def one_t(code): try: secid = ("1." if code.startswith("6") else "0.") + code url = "https://push2.eastmoney.com/api/qt/stock/get?secid=%s&fields=f57,f58,f127,f129&invt=2" % secid raw = curl(url).decode("utf-8", "ignore") j = json.loads(raw) d = j.get("data") or {} nm = d.get("f58") or "" ind = d.get("f127") or "" con = d.get("f129") or "" parts = [p for p in (con.split(",") if con else []) if p] if ind and ind not in parts: parts = [ind] + parts return (code, {"name": nm if nm != code else "", "theme": ",".join(parts[:8])}) except Exception: return (code, {"name": "", "theme": ""}) out = {} with _cf.ThreadPoolExecutor(max_workers=40) as ex: for code, v in ex.map(one_t, codes): out[code] = v self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(json.dumps(out).encode()); return if p.path == "/api/indexkline": code, beg, end = one("code", "sh000001"), one("beg"), one("end") try: import datetime as _dt1 d1 = _dt1.datetime.strptime(beg, "%Y%m%d").strftime("%Y-%m-%d") if beg and len(beg) >= 8 else "2000-01-01" d2 = _dt1.datetime.strptime(end, "%Y%m%d").strftime("%Y-%m-%d") if end and len(end) >= 8 else _dt1.datetime.now().strftime("%Y-%m-%d") url = "https://web.ifzq.gtimg.cn/appstock/app/newfqkline/get?param=%s,day,%s,%s,2000,qfq" % (code, d1, d2) raw = curl(url).decode("utf-8", "ignore") j = json.loads(raw) days = (j.get("data") or {}).get(code, {}).get("day") or [] kl = [] if days: name = {"sh000001": "上证指数", "sz399001": "深证成指", "sz399006": "创业板指", "sh000300": "沪深300"}.get(code, code) for row in days: try: d, o, c, h, l, v = row[0], row[1], row[2], row[3], row[4], row[5] except Exception: continue dd = str(d).replace("-", "") if beg and dd < beg: continue if end and dd > end: continue try: vv = str(int(float(v))) if v else "0" except Exception: vv = "0" kl.append("%s,%s,%s,%s,%s,%s,0" % (dd, o, c, h, l, vv)) data = json.dumps({"data": {"name": name, "klines": kl}}).encode() except Exception as e: data = json.dumps({"data": {"name": code, "klines": []}}).encode() self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/lhb": import concurrent.futures as _cf d = one("date") url = "https://datacenter-web.eastmoney.com/api/data/v1/get?" + urllib.parse.urlencode({ "reportName": "RPT_DAILYBILLBOARD_DETAILSNEW", "columns": "SECURITY_CODE,SECURITY_NAME_ABBR,CLOSE_PRICE,CHANGE_RATE,BILLBOARD_NET_AMT,BILLBOARD_BUY_AMT,BILLBOARD_SELL_AMT,EXPLAIN,TURNOVERRATE,TRADE_DATE", "filter": "(TRADE_DATE%3D'" + d + "')" if d else "", "pageNumber": one("pn", "1"), "pageSize": one("pz", "30"), "sortColumns": "BILLBOARD_NET_AMT", "sortTypes": "-1", }) try: data = curl(url) if b'"result"' not in data: data = b'{"result":{"data":[]}}' except Exception: data = b'{"result":{"data":[]}}' self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/hotconcept": fs = "m:90+t:3" url = EA_LIST + "?pn=1&pz=" + one("pz", "20") + "&po=1&np=1&fltt=2&invt=2&fid=f3&fs=" + fs + "&fields=f12,f14,f2,f3,f62" try: data = get(url) if b"diff" not in data: raise RuntimeError("empty") except Exception: data = curl(url) self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/name": codes = one("codes") if codes: out = {} try: cl = [c.strip() for c in codes.split(",") if c.strip()] pre = lambda c: ("sh" if c.startswith("6") else "sz") + c # 腾讯支持一次查多个,用逗号一次拉回,避免逐只curl太慢超时 t = curl(TX + ",".join(pre(c) for c in cl[:80])).decode("gbk", "ignore") for line in t.strip().splitlines(): m = line.split("~") if len(m) > 2: code = m[2] out[code] = m[1] if m[1] else code for c in cl: out.setdefault(c, c) except Exception as e: for c in codes.split(","): out.setdefault(c.strip(), c.strip()) import json as _j self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers(); self.wfile.write(_j.dumps(out).encode()); return code = one("code") pre = "sh" if code.startswith("6") else "sz" data = curl(TX + pre + code) self.send_response(200); self.send_header("Content-Type", "text/plain;charset=gbk"); self.end_headers(); self.wfile.write(data); return except Exception as e: try: self.send_response(200); self.send_header("Content-Type", "application/json"); self.end_headers() self.wfile.write(json.dumps({"error": str(e), "path": p.path}).encode()) except Exception: pass return super().do_GET() def do_POST(self): _touch() p = urllib.parse.urlparse(self.path) try: n = int(self.headers.get("Content-Length", "0")) raw = self.rfile.read(n).decode("utf-8", "ignore") if n else "" j = json.loads(raw or "{}") except Exception: j = {} try: if p.path == "/api/login": user = str(j.get("user", "")).strip() pw = str(j.get("pass", "")) cap = str(j.get("captcha", "")) token = str(j.get("token", "")) cid = str(j.get("clientId", "")).strip() or machine_id() if not user or not pw: out = {"ok": False, "msg": "请输入账号和密码"} elif not check_captcha(token, cap): out = {"ok": False, "msg": "验证码错误或已过期"} else: rc = mqtt_auth(user, pw, cid) if rc == 0: out = {"ok": True, "msg": "登录成功", "clientId": cid} elif rc in (4, 5): out = {"ok": False, "msg": "账号或密码错误(或账号无权限)"} elif rc == -1: out = {"ok": False, "msg": "无法连接MQTT服务器(%s:%d)" % (MQTT_HOST, MQTT_PORT)} else: out = {"ok": False, "msg": "登录失败(MQTT code %s)" % str(rc)} data = json.dumps(out, ensure_ascii=False).encode("utf-8") self.send_response(200); self.send_header("Content-Type", "application/json; charset=utf-8"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/upgrade": user = str(j.get("user", "")).strip() pw = str(j.get("pass", "")) try: if not user or not pw: out = {"ok": False, "msg": "缺少账号"} else: out = upgrade_check(user, pw) except RuntimeError as e: m = str(e) if m.startswith("mqtt-rc-"): out = {"ok": False, "msg": "账号或密码错误(或账号无权限)"} else: out = {"ok": False, "msg": "连接升级服务失败:%s" % m} except Exception as e: out = {"ok": False, "msg": "连接升级服务失败:%s" % str(e)} data = json.dumps(out, ensure_ascii=False).encode("utf-8") self.send_response(200); self.send_header("Content-Type", "application/json; charset=utf-8"); self.end_headers(); self.wfile.write(data); return if p.path == "/hub/publish": import base64 as _b try: if not HUB_KEY or str(j.get("key", "")) != HUB_KEY: out = {"ok": False, "msg": "无权限"} else: version = str(j.get("version", "")).strip() notes = str(j.get("notes", "")) fmap = j.get("files") or {} blobs = {} for nm, b64 in fmap.items(): blobs[str(nm)] = _b.b64decode(b64.encode()) m = hub_publish_package(version, notes, blobs) out = {"ok": True, "latest": version, "files": [e["name"] for e in m["files"]]} except Exception as e: out = {"ok": False, "msg": "发布失败:%s" % str(e)} data = json.dumps(out, ensure_ascii=False).encode("utf-8") self.send_response(200); self.send_header("Content-Type", "application/json; charset=utf-8"); self.end_headers(); self.wfile.write(data); return if p.path == "/api/upgrade/dl": user = str(j.get("user", "")).strip() pw = str(j.get("pass", "")) try: if not user or not pw: out = {"ok": False, "msg": "缺少账号"} else: out = upgrade_apply(user, pw) if out.get("applied"): import threading as _th _th.Timer(1.0, _self_restart).start() except RuntimeError as e: m = str(e) if m.startswith("mqtt-rc-"): out = {"ok": False, "msg": "账号或密码错误(或账号无权限)"} else: out = {"ok": False, "msg": "升级失败:%s" % m} except Exception as e: out = {"ok": False, "msg": "升级失败:%s" % str(e)} data = json.dumps(out, ensure_ascii=False).encode("utf-8") self.send_response(200); self.send_header("Content-Type", "application/json; charset=utf-8"); self.end_headers(); self.wfile.write(data); return except Exception as e: try: self.send_response(200); self.send_header("Content-Type", "application/json; charset=utf-8"); self.end_headers() self.wfile.write(json.dumps({"ok": False, "msg": str(e)}, ensure_ascii=False).encode("utf-8")) except Exception: pass return super().do_POST() if __name__ == "__main__": try: from http.server import ThreadingHTTPServer except ImportError: from http.server import HTTPServer from socketserver import ThreadingMixIn class ThreadingHTTPServer(ThreadingMixIn, HTTPServer): daemon_threads = True import threading as _th, sys as _sys, socket as _sk host = os.environ.get("HOST", "127.0.0.1") port = int(os.environ.get("PORT", "8000")) _probe = _sk.socket(_sk.AF_INET, _sk.SOCK_STREAM) _probe.settimeout(1) _alive = (_probe.connect_ex((host, port)) == 0) _probe.close() if _alive: print("server already running on %s:%d, exit." % (host, port)) _sys.exit(0) try: srv = ThreadingHTTPServer((host, port), H) except OSError: print("server already running on %s:%d, exit." % (host, port)) _sys.exit(0) def _watch(): if os.environ.get("HUB_MODE", "") == "1": return while True: _time.sleep(5) if HEART["seen"] and (_time.time() - HEART["beat"] > 20): print("no client heartbeat, shutting down.") try: srv.shutdown() except Exception: pass return _th.Thread(target=_watch, daemon=True).start() print("打开 http://%s:%d/index.html" % (host, port)) srv.serve_forever()