You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flask HTTP代理服务器如何支持CONNECT请求并保留日志记录功能?

Flask HTTP代理服务器如何支持CONNECT请求并保留日志记录功能?

嗨,我明白你的问题了——你的Flask代理目前没处理CONNECT请求,而这正是代理HTTPS站点时必须的核心逻辑(当你用代理访问https://开头的地址时,客户端会先发送CONNECT请求来建立加密隧道)。咱们来一步步修改代码,让它支持CONNECT方法,同时完整保留你的Firebase日志功能。

问题根源

当前你的Flask路由只声明了处理GET/POST/PUT/DELETE/PATCH/OPTIONS这些普通HTTP方法,完全没包含CONNECT;而且CONNECT的处理逻辑和普通HTTP请求天差地别——它不是要转发请求内容,而是要在客户端和目标服务器之间建立一条TCP隧道,让两者直接交换数据,这就不能用你之前的requests库转发逻辑了。

修改后的完整代码

下面是整合了CONNECT支持和日志保留的完整代码,我标注了关键改动:

from flask import Flask, request, Response
import requests
import socket
import firebase_admin
from firebase_admin import credentials, firestore
from datetime import datetime
from select import select

# === FIREBASE SETUP ===
cred = credentials.Certificate("service.json")
firebase_admin.initialize_app(cred)
db = firestore.client()
collection_name = "requests"

# === FLASK SETUP ===
app = Flask(__name__)
LOCAL_IP = socket.gethostbyname(socket.gethostname())
PORT = 5100

def log_to_firestore(entry):
    try:
        db.collection(collection_name).add(entry)
    except Exception as e:
        print(f"Firestore Error: {e}")

def handle_connect():
    # 解析CONNECT请求的目标主机和端口(格式:host:port)
    target = request.host
    if ':' not in target:
        target_host, target_port = target, 443  # 默认HTTPS端口
    else:
        target_host, target_port = target.split(':', 1)
        target_port = int(target_port)

    # 记录CONNECT请求的初始日志
    connect_entry = {
        "timestamp": datetime.utcnow().isoformat(),
        "method": "CONNECT",
        "target_host": target_host,
        "target_port": target_port,
        "client_ip": request.remote_addr,
        "status": "Initiated"
    }
    log_to_firestore(connect_entry)

    try:
        # 建立到目标服务器的TCP连接
        server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        server_socket.connect((target_host, target_port))
    except Exception as e:
        # 记录连接失败的错误日志
        error_entry = connect_entry.copy()
        error_entry["error"] = str(e)
        error_entry["status"] = "Failed"
        log_to_firestore(error_entry)
        return f"Failed to establish tunnel to {target_host}:{target_port}", 503

    # 更新日志为隧道已建立
    success_entry = connect_entry.copy()
    success_entry["status"] = "Established"
    log_to_firestore(success_entry)

    # 返回200响应告诉客户端隧道已就绪
    response = Response(status=200, headers={"Connection": "keep-alive"})

    # 定义双向数据转发的生成器函数
    def tunnel_stream():
        try:
            # 设置socket为非阻塞模式,用select实现多路复用
            server_socket.setblocking(False)
            client_stream = request.environ['wsgi.input']
            client_stream.setblocking(False)
            client_fd = client_stream.fileno()
            server_fd = server_socket.fileno()

            while True:
                # 监听客户端和服务器的可读事件
                read_ready, _, _ = select([client_fd, server_fd], [], [], 5)
                
                if client_fd in read_ready:
                    # 从客户端读数据,转发到服务器
                    client_data = client_stream.read(4096)
                    if not client_data:
                        break
                    server_socket.sendall(client_data)
                
                if server_fd in read_ready:
                    # 从服务器读数据,转发到客户端
                    server_data = server_socket.recv(4096)
                    if not server_data:
                        break
                    yield server_data
        except Exception as e:
            # 转发过程中出现错误,记录日志
            error_entry = connect_entry.copy()
            error_entry["error"] = str(e)
            error_entry["status"] = "Error during tunnel"
            log_to_firestore(error_entry)
        finally:
            # 关闭连接,记录隧道关闭日志
            server_socket.close()
            close_entry = connect_entry.copy()
            close_entry["status"] = "Closed"
            close_entry["end_timestamp"] = datetime.utcnow().isoformat()
            log_to_firestore(close_entry)

    # 将转发函数作为响应内容返回
    response.response = tunnel_stream()
    return response

@app.route('/', defaults={'path': ''}, methods=['GET', 'POST', 'PUT', 'DELETE', 'PATCH', 'OPTIONS', 'CONNECT'])
@app.route('/<path:path>', methods=['GET', 'POST', 'PUT', 'DELETE', 'PATCH', 'OPTIONS', 'CONNECT'])
def proxy(path):
    # 优先处理CONNECT请求
    if request.method == 'CONNECT':
        return handle_connect()

    # 原来的普通HTTP请求转发逻辑(完全保留原有功能)
    if request.url.startswith('http'):
        target_url = request.url
    else:
        target_url = request.args.get('url')
        if not target_url:
            if path and '://' in path:
                target_url = path
            else:
                return "Missing target URL", 400

    try:
        request_data = request.get_data()
        headers = {k: v for k, v in request.headers.items() if k.lower() != 'host'}

        # Forward the request
        response = requests.request(
            method=request.method,
            url=target_url,
            headers=headers,
            data=request_data,
            cookies=request.cookies,
            allow_redirects=False
        )

        # Prepare Firestore log entry
        entry = {
            "timestamp": datetime.utcnow().isoformat(),
            "method": request.method,
            "target_url": target_url,
            "client_ip": request.remote_addr,
            "request": {
                "headers": dict(request.headers),
                "body": request_data.decode(errors='ignore'),
            },
            "response": {
                "status": response.status_code,
                "headers": dict(response.headers),
                "body": response.text[:1000]  # Limit for Firestore size
            }
        }

        log_to_firestore(entry)

        return Response(response.content, 
                        status=response.status_code, 
                        headers=dict(response.headers))

    except Exception as e:
        error_entry = {
            "timestamp": datetime.utcnow().isoformat(),
            "method": request.method,
            "target_url": target_url,
            "client_ip": request.remote_addr,
            "error": str(e)
        }
        log_to_firestore(error_entry)
        return f"Error: {str(e)}", 500

if __name__ == '__main__':
    print(f"Proxy Server running at: http://{LOCAL_IP}:{PORT}")
    # 注意:生产环境请关闭debug模式,避免干扰socket转发
    app.run(host='0.0.0.0', port=PORT, debug=False)

关键改动说明

  1. 添加CONNECT方法支持:在Flask路由的methods参数里新增了'CONNECT',这样Flask不会再返回405错误。
  2. 单独处理CONNECT请求:新增了handle_connect函数,专门处理隧道建立逻辑:
    • 解析目标主机和端口,建立TCP连接到目标服务器
    • 用select实现客户端和服务器之间的双向数据转发(非阻塞IO,避免卡死)
  3. 完整的CONNECT日志:记录了请求的发起、建立成功、错误、关闭等全生命周期状态,和普通请求的日志格式统一,方便在Firebase中查询。
  4. 兼容原有逻辑:普通HTTP请求的转发和日志逻辑完全保留,不影响你原来的使用场景。

测试验证

修改代码后,你可以用以下命令测试CONNECT请求(HTTPS代理场景):

curl -x http://22.22.5.30:5100 https://api.ipify.org

同时你可以检查Firebase的requests集合,应该能看到这条CONNECT请求的完整日志记录。

注意事项

  • 请确保你的Cloud Run实例端口配置和代码中的PORT一致,防火墙规则放行对应端口。
  • Flask的debug模式会启用重载器,可能干扰socket转发,生产环境务必关闭debug=False。
  • 如果你的代理需要处理大流量,可能需要优化socket转发的缓冲区大小,或者增加超时时间。

备注:内容来源于stack exchange,提问作者Ahmed Zaidan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.13 19:48:08