модуль ActionCable::Channel::Streams
Потоки канала Action Cable
Streams позволяют каналам направлять трансляции подписчикам. Трансляция, как обсуждалось ранее, это очередь pubsub, где любые данные, помещённые в неё, автоматически отправляются клиентам, которые подключены в данный момент. Однако это чисто онлайн-очередь. Если вы не ведёте трансляцию трансляции в тот момент, когда она отправляет обновление, вы не получите это обновление, даже если подключитесь после того, как оно было отправлено.
Чаще всего транслируемая трансляция отправляется непосредственно подписчику на стороне клиента. Канал просто выступает в роли посредника между двумя сторонами (транслятором и подписчиком канала). Вот пример канала, который позволяет подписчикам получать все новые комментарии на данной странице:
class CommentsChannel < ApplicationCable::Channel
def follow(data)
stream_from "comments_for_#{data['recording_id']}"
end
def unfollow
stop_all_streams
end
end
На основе приведенного выше примера подписчики этого канала получат любые данные, которые будут помещены в, скажем, comments_for_45 трансляцию, как только они будут помещены.
Пример трансляции для этого канала выглядит следующим образом:
ActionCable.server.broadcast "comments_for_45", { author: 'DHH', content: 'Rails is just swell' }
Если у вас есть поток, связанный с моделью, то используемая трансляция может быть сгенерирована из модели и канала. Следующий пример подпишется на трансляцию, такую как comments:Z2lkOi8vVGVzdEFwcC9Qb3N0LzE.
class CommentsChannel < ApplicationCable::Channel
def subscribed
post = Post.find(params[:id])
stream_for post
end
end
Затем вы можете транслировать в этот канал, используя:
CommentsChannel.broadcast_to(@post, @comment)
Если вам не нужно просто транслировать трансляцию без фильтрации подписчику, вы также можете указать обратный вызов, который позволит вам изменить то, что отправляется. Приведённый ниже пример показывает, как вы можете использовать это для предоставления интроспекции производительности в процессе:
class ChatChannel < ApplicationCable::Channel
def subscribed
@room = Chat::Room[params[:room_number]]
stream_for @room, coder: ActiveSupport::JSON do |message|
if message['originated_at'].present?
elapsed_time = (Time.now.to_f - message['originated_at']).round(2)
ActiveSupport::Notifications.instrument :performance, measurement: 'Chat.message_delay', value: elapsed_time, action: :timing
logger.info "Message took #{elapsed_time}s to arrive"
end
transmit message
end
end
end
Вы можете остановить трансляцию всех трансляций, вызвав stop_all_streams.
Методы публичного экземпляра
# File actioncable/lib/action_cable/channel/streams.rb, line 122
def stop_all_streams
streams.each do |broadcasting, callback|
pubsub.unsubscribe broadcasting, callback
logger.info "#{self.class.name} stopped streaming from #{broadcasting}"
end.clear
end Отписывает все потоки, связанные с этим каналом, от очереди pubsub.
# File actioncable/lib/action_cable/channel/streams.rb, line 117 def stop_stream_for(model) stop_stream_from(broadcasting_for(model)) end
Отписывает потоки для model.
# File actioncable/lib/action_cable/channel/streams.rb, line 108
def stop_stream_from(broadcasting)
callback = streams.delete(broadcasting)
if callback
pubsub.unsubscribe(broadcasting, callback)
logger.info "#{self.class.name} stopped streaming from #{broadcasting}"
end
end Отписывает потоки от указанной broadcasting.
# File actioncable/lib/action_cable/channel/streams.rb, line 103 def stream_for(model, callback = nil, coder: nil, &block) stream_from(broadcasting_for(model), callback || block, coder: coder) end
Начать трансляцию очереди pubsub для model в этом канале. По желанию, вы можете передать callback, которое будет использоваться вместо по умолчанию, просто передавая обновления напрямую подписчику.
Передайте coder: ActiveSupport::JSON для декодирования сообщений в формате JSON перед передачей в обратный вызов. По умолчанию coder: nil не выполняет декодирования, передает необработанные сообщения.
# File actioncable/lib/action_cable/channel/streams.rb, line 78
def stream_from(broadcasting, callback = nil, coder: nil, &block)
broadcasting = String(broadcasting)
# Don't send the confirmation until pubsub#subscribe is successful
defer_subscription_confirmation!
# Build a stream handler by wrapping the user-provided callback with
# a decoder or defaulting to a JSON-decoding retransmitter.
handler = worker_pool_stream_handler(broadcasting, callback || block, coder: coder)
streams[broadcasting] = handler
connection.server.event_loop.post do
pubsub.subscribe(broadcasting, handler, lambda do
ensure_confirmation_sent
logger.info "#{self.class.name} is streaming from #{broadcasting}"
end)
end
end Начать трансляцию с указанной broadcasting очереди pubsub. По желанию, вы можете передать callback, которое будет использоваться вместо по умолчанию, просто передавая обновления напрямую подписчику. Передайте coder: ActiveSupport::JSON для декодирования сообщений в формате JSON перед передачей в обратный вызов. По умолчанию coder: nil не выполняет декодирования, передает необработанные сообщения.
# File actioncable/lib/action_cable/channel/streams.rb, line 131
def stream_or_reject_for(model)
if model
stream_for model
else
reject
end
end Вызывает stream_for с заданной model, если она присутствует, чтобы начать трансляцию, в противном случае отклоняет подписку.
© 2004–2021 David Heinemeier Hansson
Licensed under the MIT License.