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

Cloudflare Durable Object中WebSocket Hibernation API实现问题排查

迁移Cloudflare Durable Object到WebSocket Hibernation API时的客户端状态丢失问题

我基于Cloudflare Durable Objects和WebSocket开发了Worker服务,由于WebSocket运行成本较高,计划切换到官方的WebSocket Hibernation API。原本以为只需移除原有的事件监听器,改用内置的webSocketMessage和webSocketClose方法即可完成迁移,但实际运行后发现代码存在问题:迁移后的代码无法正确保存this.clients客户端连接列表以及programConnection对象,导致向程序发送事件的核心逻辑无法触发。原标准WebSocket版本的代码运行完全正常,我查阅了所有相关文档,但可用资料非常有限。

原标准WebSocket API代码

export class DurableObject_Manager {
  constructor(state, env) {
    this.state = state;
    this.env = env;
    this.clients = new Set(); // Set of active WebSocket connections
    this.programConnection = null;
  }

  async fetch(request) {
    const upgradeHeader = request.headers.get("Upgrade");

    if (upgradeHeader === "websocket") {
      const [client, server] = Object.values(new WebSocketPair());
      server.accept();

      // Checks if the initial connection is the program
      if (!this.programConnection) {
        this.programConnection = server;
      } else {
        if (this.programConnection.readyState === WebSocket.OPEN) {
          this.programConnection.send(
            JSON.stringify({ Evento: "sincronizar_fila", Comando: "" })
          );
        }
      }

      // Adds the connection to the set
      this.clients.add(server);

      // When connecting, sends stored songs in Durable Object, if any exist
      const storedMusics = await this.state.storage.get("musicas");
      if (storedMusics) {
        server.send(JSON.stringify({ Evento: "musicas", Comando: storedMusics }));
      }

      server.addEventListener("message", async (data) => {
        let message;

        try {
          message = JSON.parse(data.data);
        } catch (error) {
          console.error("Invalid msg received", data.data);
          return;
        }

        // verify if the message event is "musicas"
        if (message.Evento === "musicas") {
          const existingMusics = await this.state.storage.get("musicas");
          if (!existingMusics) {
            // Saves to storage only if there are no songs saved
            await this.state.storage.put("musicas", message.Comando);
          }
        } 

          this.broadcast(data.data, server);

      });

      server.addEventListener("close", async () => {
        this.clients.delete(server);

        if (server === this.programConnection) {
          console.log("Main program disconnected. Closing all connections...");
          this.programConnection = null;

          for (const client of this.clients) {
            client.close(1000, "Main program disconnected");
          }

          this.clients.clear();
          await this.state.storage.delete("musicas");
        }

        if (this.clients.size === 0) {
          console.log("All clients disconnected. Clearing storage...");
          await this.state.storage.delete("musicas");
        }
      });

      return new Response(null, { status: 101, webSocket: client });
    }

    return new Response("The request is not a websocket.", { status: 400 });
  }

  broadcast(data, sender) {
    // Send the command to all connected clients except the sender
    for (const client of this.clients) {
      if (client !== sender && client.readyState === WebSocket.OPEN) {
        client.send(data);
      }
    }
  }
}

迁移后的WebSocket Hibernation API尝试代码

export class DurableObject_Manager {
  constructor(state, env) {
    this.state = state;
    this.env = env;
    this.clients = new Set(); // Set of active WebSocket connections
    this.programConnection = null;
  }

  async fetch(request) {
    const upgradeHeader = request.headers.get("Upgrade");

    if (upgradeHeader === "websocket") {
      const [client, server] = Object.values(new WebSocketPair());

      this.state.acceptWebSocket(server);

      // Checks if the initial connection is the program
      if (!this.programConnection) {
        this.programConnection = server;
      } else {
        if (this.programConnection.readyState === WebSocket.OPEN) {
          this.programConnection.send(
            JSON.stringify({ Evento: "sincronizar_fila", Comando: "" })
          );
        }
      }

      // Add the connection to the set
      this.clients.add(server);

      // When connecting, send the stored songs, if any exists
      const storedMusics = await this.state.storage.get("musicas");
      if (storedMusics) {
        server.send(JSON.stringify({ Evento: "musicas", Comando: storedMusics }));
      }
      return new Response(null, { status: 101, webSocket: client });
    }

    return new Response("The request is not a websocket.", { status: 400 });
  }

  async webSocketMessage(ws, message) {
    let data;

    try {
      data = JSON.parse(message);
    } catch (error) {
      console.error("Invalid msg received:", message);
      return;
    }

    // Check if the event is about music
  if (data.Evento === "musicas") {
    const existingMusics = await this.state.storage.get("musicas");
    if (!existingMusics) {
      // Save to storage only if there are no songs saved
      await this.state.storage.put("musicas", data.Comando);
    }
  }
    this.broadcast(data, ws);
  }

  async webSocketClose(ws){
    console.log(ws)
    ws.close();
  }

broadcast(data, sender) {
 // Send the command to all connected clients
  for (const client of this.clients) {
    if (client !== sender && client.readyState === WebSocket.OPEN) {
      client.send(data);
    }
  }
}
}

我移除了原有的消息和关闭事件监听器,将逻辑迁移到webSocketMessage()和webSocketClose()方法中后,代码虽能运行,但this.clients客户端列表和programConnection的状态无法被正确保留,导致向程序发送事件的逻辑无法触发,寻求有相关开发经验的人士协助解决。

内容的提问来源于stack exchange,提问作者Vítor Souza

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 04:24:50