【问题标题】:BroadcastHub filtering based on "resource" the connected client is working on?BroadcastHub 过滤基于连接的客户端正在处理的“资源”?
【发布时间】:2017-10-01 08:31:30
【问题描述】:

我正在编写一个纯 websocket web 应用程序,这意味着在 websocket 升级之前没有用户/客户端步骤,更具体地说: 身份验证请求和其他通信一样通过 websockets

有/是:

  • /api/ws 上只有一个 websocket 端点
  • 连接到该端点的多个客户端
  • 多个客户的多个项目

现在,并不是每个客户端都可以访问每个项目 - 访问控制是在服务器端 (ofc) 实现的,与 websockets 本身无关。

我的问题是,我想允许协作,这意味着 N 个客户可以一起处理 1 个项目。

现在,如果其中一个客户修改了某些内容,我想通知正在从事该项目的所有其他客户。

这一点尤其重要,因为 atm 我是唯一一个在这方面工作并对其进行测试的人,这是我这边的主要疏忽,因为现在:

如果客户 A 连接到项目 X 并且客户 B 连接到项目 Y,如果其中任何一个更新了各自项目中的某些内容,则另一个会收到这些更改的通知。

现在我的 WebsocketController 比较简单,基本有这个:

private val fanIn = MergeHub.source[AllowedWSMessage].to(sink).run()
private val fanOut = source.toMat(BroadcastHub.sink[AllowedWSMessage])(Keep.right).run()

def handle: WebSocket = WebSocket.accept[AllowedWSMessage, AllowedWSMessage]
{
  _ => Flow.fromSinkAndSource(fanIn, fanOut)
}

现在根据我的理解,我需要的是

1) 每个项目有多个 websocket 端点,例如 /api/{project_identifier}/ws

(X)OR

2) 根据他们正在工作的项目来拆分 WebSocket 连接/连接的客户端的一些方法。

因为我不想走路线 1) 我将分享我对 2) 的想法:

我目前看不到解决方法的问题是,我可以轻松地在服务器端创建一些集合,在其中存储在任何给定时刻哪个用户正在处理哪个项目(例如,如果他们选择/切换一个项目,客户端将其发送到服务器并存储此信息)

但是我仍然有那个fanOut,所以这不能解决我关于 WebSocket/AkkaStreams 的问题。

BroadcastHub 上是否有一些魔法(过滤)可以满足我的要求?

编辑:在尝试但未能应用@James Roper 的良好提示之后,现在在这里分享我的整个 websocket 逻辑:

 class WebSocketController @Inject()(implicit cc: ControllerComponents, ec: ExecutionContext, system: ActorSystem, mat: Materializer) extends AbstractController(cc)

{ val logger: Logger = Logger(this.getClass())

type WebSocketMessage = Array[Byte]

import scala.concurrent.duration._

val tickingSource: Source[WebSocketMessage, Cancellable] =
  Source.tick(initialDelay = 1 second, interval = 10 seconds, tick = NotUsed)
    .map(_ => Wrapper().withKeepAlive(KeepAlive()).toByteArray)

private val generalActor = system.actorOf(Props
{
  new myActor(system, "generalActor")
}, "generalActor")

private val serverMessageSource = Source
  .queue[WebSocketMessage](10, OverflowStrategy.backpressure)
  .mapMaterializedValue
  { queue => generalActor ! InitTunnel(queue) }

private val sink: Sink[WebSocketMessage, NotUsed] = Sink.actorRefWithAck(generalActor, InternalMessages.Init(), InternalMessages.Acknowledged(), InternalMessages.Completed())
private val source: Source[WebSocketMessage, Cancellable] = tickingSource.merge(serverMessageSource)

private val fanIn = MergeHub.source[WebSocketMessage].to(sink).run()
private val fanOut = source.toMat(BroadcastHub.sink[WebSocketMessage])(Keep.right).run()

// TODO switch to WebSocket.acceptOrResult
def handle: WebSocket = WebSocket.accept[WebSocketMessage, WebSocketMessage]
  {
    //_ => createFlow()
    _ => Flow.fromSinkAndSource(fanIn, fanOut)
  }

private val projectHubs = TrieMap.empty[String, (Sink[WebSocketMessage, NotUsed], Source[WebSocketMessage, NotUsed])]

private def buildProjectHub(projectName: String) =
{
  logger.info(s"building projectHub for $projectName")

  val projectActor = system.actorOf(Props
  {
    new myActor(system, s"${projectName}Actor")
  }, s"${projectName}Actor")

  val projectServerMessageSource = Source
    .queue[WebSocketMessage](10, OverflowStrategy.backpressure)
    .mapMaterializedValue
    { queue => projectActor ! InitTunnel(queue) }

  val projectSink: Sink[WebSocketMessage, NotUsed] = Sink.actorRefWithAck(projectActor, InternalMessages.Init(), InternalMessages.Acknowledged(), InternalMessages.Completed())
  val projectSource: Source[WebSocketMessage, Cancellable] = tickingSource.merge(projectServerMessageSource)

  val projectFanIn = MergeHub.source[WebSocketMessage].to(projectSink).run()
  val projectFanOut = projectSource.toMat(BroadcastHub.sink[WebSocketMessage])(Keep.right).run()

  (projectFanIn, projectFanOut)
}

private def getProjectHub(userName: String, projectName: String): Flow[WebSocketMessage, WebSocketMessage, NotUsed] =
{
  logger.info(s"trying to get projectHub for $projectName")

  val (sink, source) = projectHubs.getOrElseUpdate(projectName, {
    buildProjectHub(projectName)
  })

  Flow.fromSinkAndSourceCoupled(sink, source)
}

private def extractUserAndProject(msg: WebSocketMessage): (String, String) =
{
  Wrapper.parseFrom(msg).`type` match
  {
    case m: MessageType =>
      val message = m.value
      (message.userName, message.projectName)
    case _ => ("", "")
  }
}

private def createFlow(): Flow[WebSocketMessage, WebSocketMessage, NotUsed] =
{
  // broadcast source and sink for demux/muxing multiple chat rooms in this one flow
  // They'll be provided later when we materialize the flow
  var broadcastSource: Source[WebSocketMessage, NotUsed] = null
  var mergeSink: Sink[WebSocketMessage, NotUsed] = null

  Flow[WebSocketMessage].map
  {
    m: WebSocketMessage =>
    val msg = Wrapper.parseFrom(m)
    logger.warn(s"client sent project related message: ${msg.toString}");
    m
  }.map
    {
      case isProjectRelated if !extractUserAndProject(isProjectRelated)._2.isEmpty =>
        val (userName, projectName) = extractUserAndProject(isProjectRelated)

        logger.info(s"userName: $userName, projectName: $projectName")
        val projectFlow = getProjectHub(userName, projectName)

        broadcastSource.filter
        {
          msg =>
            val (_, project) = extractUserAndProject(msg)
            logger.info(s"$project == $projectName")
            (project == projectName)
        }
          .via(projectFlow)
          .runWith(mergeSink)

        isProjectRelated

      case other =>
      {
        logger.info("other")
        other
      }
    } via {
      Flow.fromSinkAndSourceCoupledMat(BroadcastHub.sink[WebSocketMessage], MergeHub.source[WebSocketMessage])
      {
        (source, sink) =>
          broadcastSource = source
          mergeSink = sink

          source.filter(extractUserAndProject(_)._2.isEmpty)
            .map
            { x => logger.info("Non project related stuff"); x }
            .via(Flow.fromSinkAndSource(fanIn, fanOut))
            .runWith(sink)

          NotUsed
      }
    }
}

}

解决方案/想法我是如何理解的:

1) 我们有一个“包装流程”,其中我们有一个为 null 的 broadcastSource 和 mergeSink,直到我们在外部 } via { 块中实现它们

2) 在那个“包装流程”中,我们映射每个元素来检查它。

I) 如果是项目相关的,我们

a) 为项目获取/创建自己的子流程 b) 根据项目名称过滤元素 c)让通过过滤器的那些被子/项目流消耗,以便连接到项目的每个人都得到那个元素

II) 如果它与项目无关,我们只是传递它

3) 我们的包装流程通过“按需”具体化流程进行,在具体化它的via 中,我们让与项目无关的元素分发到所有连接的 Web 套接字客户端。

总结一下:我们有一个用于 websocket 连接的“包装流”,它可以通过 projectFlow 或 generalFlow 进行,具体取决于它正在处理的消息/元素。

我现在的问题是(这似乎是微不足道的,但我不知何故在挣扎)每条消息都应该进入myActor(atm),并且也应该有消息从那里发出(见serverMesssageSourcesource)

但上面的代码正在创建不确定的结果,例如一个客户端发送 2 条消息,但有 4 条正在处理(根据服务器发回的日志和结果),有时消息在从控制器到参与者的途中突然丢失。

我无法解释,但如果我只留下_ => Flow.fromSinkAndSource(fanIn, fanOut),每个人都会得到一切,但至少如果只有一个客户,它会完全符合预期(显然:))

【问题讨论】:

    标签: playframework websocket akka broadcast akka-stream


    【解决方案1】:

    我实际上建议使用 Play 的socket.io support。这提供了命名空间,从你的描述中我可以看出,它可以直接实现你想要的——每个命名空间都是它自己独立管理的流,但所有命名空间都使用同一个 WebSocket。我wrote a blog post 谈谈你今天为什么会选择使用 socket.io。

    如果您不想使用 socket.io,我在这里有一个示例(它使用 socket.io,但不使用 socket.io 命名空间,因此可以很容易地适应直接在 WebSockets 上运行),它显示多聊天室协议 - 它将消息馈送到 BroadcastHub,然后用户当前所在的每个聊天室都有一个订阅中心(对您而言,每个项目都有一个订阅)。这些订阅中的每一个都会过滤来自中心的消息,以仅包含该订阅聊天室的消息,然后将消息馈送到该聊天室 MergeHub。

    这里突出显示的代码根本不是socket.io特有的,如果你可以将WebSocket连接适配为ChatEvent的流,你可以这样使用:

    https://github.com/playframework/play-socket.io/blob/c113e74a4d9b435814df1ccdc885029c397d9179/samples/scala/multi-room-chat/app/chat/ChatEngine.scala#L84-L125

    为了满足您通过每个人都连接到的广播频道引导非项目特定消息的要求,首先,创建该频道:

    val generalFlow = {
      val (sink, source) = MergeHub.source[NonProjectSpecificEvent]
        .toMat(BroadcastHub.sink[NonProjectSpecificEvent])(Keep.both).run
      Flow.fromSinkAndSourceCoupled(sink, source)
    }
    

    然后,当每个连接的 WebSocket 的广播接收器/源连接时,附加它(这是来自聊天示例:

    } via {
      Flow.fromSinkAndSourceCoupledMat(BroadcastHub.sink[YourEvent], MergeHub.source[YourEvent]) { (source, sink) =>
        broadcastSource = source
        mergeSink = sink
    
        source.filter(_.isInstanceOf[NonProjectSpecificEvent])
          .via(generalFlow)
          .runWith(sink)
    
        NotUsed
      }
    }
    

    【讨论】:

    • 所以对我来说有点复杂,因为我没有一路 JSON,而是 protobuf,所以我需要调查每个请求并检查用户是否加入项目,离开项目等。
    • 如果不是 JSON 也没关系,如果是 protobuf,您可以使用 protobuf 将流从字节数组转换为对象。
    • 理论上是的。问题是 websocket.accept 使用 requestheader,我不知道如何到达这里的正文。到现在为止,我在已经获得字节的演员后面做到了这一点。那里很容易。所以这就是我现在卡在 play 的 API 上的地方。
    • 发出 WebSocket 请求时没有请求正文。请求标头与特定的 WebSocket 标头一起发送,响应标头与特定的标头一起发送回,然后添加 WebSocket 消息开始来回发送。
    • 啊,这完全有道理。那么如何根据 websocket 内容决定构建哪个流呢?我在这里看到了一个鸡蛋问题。需要再次检查您是如何解决的。也许我需要另一个参与者,一个接收器,每个请求都进入其中,并通过正确的流程进行委托。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-11-05
    • 1970-01-01
    • 2013-09-26
    • 1970-01-01
    • 1970-01-01
    • 2018-04-23
    相关资源
    最近更新 更多