from fastapi import FastAPI, Request, Response from fastapi.responses import JSONResponse, PlainTextResponse import httpx import os import sys from collections import deque app = FastAPI() # 内存日志队列(最多保存1000条) log_buffer = deque(maxlen=1000) def log_message(msg): """记录日志到内存和标准输出""" log_buffer.append(msg) print(msg, flush=True) @app.get("/") @app.head("/") async def root(): return {"status": "ok", "service": "HTTP Proxy Node"} @app.get("/logs") async def get_logs(): """获取最近的日志""" return PlainTextResponse("\n".join(log_buffer)) @app.post("/api/proxy") async def proxy(request: Request): """标准 HTTP 代理 API""" import time start_time = time.time() try: # 解析请求 parse_start = time.time() data = await request.json() target_url = data.get("url") method = data.get("method", "GET") headers = data.get("headers", {}) body = data.get("body") log_message(f"[PROXY] {method} {target_url[:80]} - parse: {(time.time()-parse_start)*1000:.0f}ms") if not target_url: return JSONResponse( {"error": "Missing url parameter"}, status_code=400 ) # 转换 body(如果是 base64) if body and isinstance(body, str): import base64 try: body = base64.b64decode(body) except: body = body.encode('utf-8') # 发送请求(异步) fetch_start = time.time() log_message(f"[FETCH_START] {method} {target_url[:80]}") try: # 设置更详细的超时 timeout_config = httpx.Timeout( connect=10.0, # 连接超时 read=30.0, # 读取超时 write=10.0, # 写入超时 pool=10.0 # 连接池超时 ) async with httpx.AsyncClient(timeout=timeout_config, follow_redirects=True) as client: resp = await client.request( method=method, url=target_url, headers=headers, content=body ) log_message(f"[FETCH] {resp.status_code} {len(resp.content)} bytes - {(time.time()-fetch_start)*1000:.0f}ms") except httpx.TimeoutException as e: error_msg = f"{type(e).__name__}: {str(e) or repr(e)}" log_message(f"[FETCH_TIMEOUT] {error_msg} - {(time.time()-fetch_start)*1000:.0f}ms") raise except httpx.HTTPError as e: error_msg = f"{type(e).__name__}: {str(e) or repr(e)}" log_message(f"[FETCH_HTTP_ERROR] {error_msg} - {(time.time()-fetch_start)*1000:.0f}ms") raise except Exception as fetch_error: error_msg = f"{type(fetch_error).__name__}: {str(fetch_error) or repr(fetch_error)}" log_message(f"[FETCH_ERROR] {error_msg} - {(time.time()-fetch_start)*1000:.0f}ms") import traceback log_message(f"[TRACEBACK] {traceback.format_exc()}") raise # 返回响应(过滤掉会冲突的 headers) response_headers = {} for k, v in resp.headers.items(): k_lower = k.lower() # 跳过这些 headers,让 FastAPI 自动处理 if k_lower not in ['content-length', 'content-encoding', 'transfer-encoding']: response_headers[k] = v log_message(f"[DONE] Total: {(time.time()-start_time)*1000:.0f}ms") return Response( content=resp.content, status_code=resp.status_code, headers=response_headers ) except Exception as e: log_message(f"[ERROR] {str(e)} - {(time.time()-start_time)*1000:.0f}ms") return JSONResponse( {"error": str(e)}, status_code=500 ) if __name__ == "__main__": import uvicorn import threading import time # 启动时自动注册 def auto_register(): time.sleep(10) # 等待服务启动 try: import httpx # 从 HF 环境变量获取 Space 信息 space_id = os.environ.get("SPACE_ID", "") # pyrq/ws-test-xxx if space_id: space_url = f"https://{space_id.replace('/', '-')}.hf.space" # 注册到 Worker(新格式) with httpx.Client(timeout=10.0) as client: client.post( "https://proxy-worker.busitest135.workers.dev/api/instances", json={ "nodeId": space_id, "url": space_url, "platform": "hf", "tags": ["auto"] } ) log_message(f"[AUTO_REGISTER] Registered {space_id} -> {space_url}") except Exception as e: log_message(f"[AUTO_REGISTER] Failed: {e}") threading.Thread(target=auto_register, daemon=True).start() port = int(os.environ.get("PORT", 7860)) uvicorn.run(app, host="0.0.0.0", port=port)