модуль ActionCable::Channel::Streams
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 120
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 115 def stop_stream_for(model) stop_stream_from(broadcasting_for(model)) end
Отписывает потоки для model.
# File actioncable/lib/action_cable/channel/streams.rb, line 106
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 101 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 76
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(record)
if record
stream_for record
else
reject
end
end Вызывает stream_for, если запись присутствует, в противном случае вызывает reject. Этот метод предназначен для вызова, когда вы ищете запись на основе параметра. Если запись найдена, запускается трансляция. Если запись равна nil, соединение отклоняется.
© 2004–2020 David Heinemeier Hansson
Licensed under the MIT License.