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

WebSocket无提示自动关闭问题:推送客户端余额更新时发送accepted消息后无法传输数据

问题分析与解决方案

看起来你的WebSocket连接在发送accepted消息后就自动关闭,核心原因是FastAPI的WebSocket端点函数执行完毕后,会自动终止连接。当前你的api_websocket_wait_update函数在发送accepted后没有任何阻塞逻辑,函数直接结束,导致连接被关闭。此外还有几处代码细节问题需要修正,我来一步步帮你解决:

核心问题与修正步骤

1. 让WebSocket端点保持运行

FastAPI的WebSocket端点需要函数持续运行才能维持连接。我们可以用asyncio.Event来等待余额更新事件,直到回调触发或者客户端主动断开。

2. 修正Token不存在的错误处理

当没有获取到access_token时,返回JSONResponse是无效的——因为WebSocket已经建立了双向连接,应该发送错误消息后主动关闭连接。

3. 修复Firebase监听器的回调参数错误

你的get_contract_listener中,on_snapshot的回调参数处理有误:单个文档的on_snapshot回调第一个参数就是DocumentSnapshot对象,不需要取doc[0],这会导致无法正确获取余额数据。

4. 正确处理异步任务提交

在同步的on_snapshot回调中执行异步函数,不能用asyncio.run(会创建新事件循环),应该用asyncio.create_task将异步任务提交到当前FastAPI的事件循环中。

修改后的完整代码

后端Python代码

import asyncio
from fastapi import WebSocket, WebSocketDisconnect
from google.cloud.firestore import DocumentSnapshot, Watch, FieldFilter
from typing import Callable, Awaitable, cast

# 修正get_contract_listener函数
def get_contract_listener(access_token: str, callback: Callable[[int | None], Awaitable[None]]) -> None | Watch:
    logger.debug('get contract listener token=%s', access_token)
    contract = db.collection('Contract').where(filter=FieldFilter('access_token', '==', access_token)).get()
    if not contract:
        logger.debug('contract not found')
        return
    async def wrap(document: DocumentSnapshot):
        await callback((document.to_dict() or {}).get('last_saldo'))
    # 修正回调参数,用create_task提交异步任务
    return cast(DocumentReference, contract[0].reference).on_snapshot(
        lambda doc, _, __: asyncio.create_task(wrap(doc))
    )

# 修正WebSocket端点函数
@router.websocket('/wait_update')
async def api_websocket_wait_update(socket: WebSocket):
    await socket.accept()
    token = socket.cookies.get('access_token')
    
    # 修正Token不存在的处理逻辑
    if token is None:
        await socket.send_json({'status': 'fail', 'detail': 'no access token'})
        await socket.close(1008)  # 1008表示权限相关的关闭码
        return
    
    update_event = asyncio.Event()
    listener = None
    
    async def callback(balance: int | None) -> None:
        nonlocal listener
        # 先取消监听器,避免重复触发
        if listener is not None:
            listener.unsubscribe()
        # 发送更新消息
        await socket.send_json({'message': 'updated', 'balance': balance})
        # 触发事件,让主函数退出,关闭连接
        update_event.set()
    
    listener = get_contract_listener(token, callback)
    if listener is None:
        await socket.send_json({'message': 'failed', 'detail': 'contract not found'})
        await socket.close(1007)
        return
    
    await socket.send_json({'message': 'accepted'})
    
    try:
        # 等待更新事件或者客户端断开连接
        await asyncio.wait_for(update_event.wait(), timeout=None)
    except WebSocketDisconnect:
        # 客户端主动断开时,清理监听器
        if listener is not None:
            listener.unsubscribe()
        logger.debug('Client disconnected')
    finally:
        # 确保连接关闭
        await socket.close()

前端JS代码(添加错误提示优化)

function create_socket() {
    const socket = new WebSocket('/api/contract/wait_update');
    socket.addEventListener('open', () => {
        console.log('open websocket connection');
    });
    socket.addEventListener('message', (e) => {
        let data;
        try {
            data = JSON.parse(e.data);
        } catch (err) {
            console.error('Failed to parse websocket message:', err);
            return;
        }
        if (data.message == 'failed') {
            console.error(`error from websocket: ${data.detail}`);
            socket.close();
            return;
        }
        console.log(`message from websocket: ${e.data}`);
        if (data.message == 'updated') {
            console.log('balance updated');
            // 执行你的余额更新逻辑...
            socket.close(); // 更新完成后主动关闭连接
        }
    });
    socket.addEventListener('close', (event) => {
        console.log(`close websocket connection, code: ${event.code}, reason: ${event.reason}`);
    });
    socket.addEventListener('error', (err) => {
        console.error('WebSocket error:', err);
    });
}

关键修改说明

  • 用asyncio.Event维持连接:端点函数会一直等待update_event被触发(余额更新时),避免函数提前结束导致连接关闭。
  • 修正Firebase监听器参数:将doc[0]改为doc,确保正确传递DocumentSnapshot对象到回调函数。
  • 异步任务提交:用asyncio.create_task替代asyncio.run,避免在现有事件循环中创建新循环引发冲突。
  • 完善错误处理:Token不存在时发送WebSocket错误消息并关闭连接,客户端断开时清理Firebase监听器,避免资源泄漏。

内容的提问来源于stack exchange,提问作者firedotguy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.27 10:12:33