forked from cybrespace/mastodon
		
	Merge branch 'krainboltgreene-broadcast-to-worker'
This commit is contained in:
		
						commit
						d13d169922
					
				
					 6 changed files with 38 additions and 19 deletions
				
			
		| 
						 | 
					@ -34,12 +34,7 @@ class FeedManager
 | 
				
			||||||
      trim(timeline_type, account.id)
 | 
					      trim(timeline_type, account.id)
 | 
				
			||||||
    end
 | 
					    end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
    broadcast(account.id, event: 'update', payload: inline_render(account, 'api/v1/statuses/show', status))
 | 
					    PushUpdateWorker.perform_async(account.id, status.id)
 | 
				
			||||||
  end
 | 
					 | 
				
			||||||
 | 
					 | 
				
			||||||
  def broadcast(timeline_id, options = {})
 | 
					 | 
				
			||||||
    options[:queued_at] = (Time.now.to_f * 1000.0).to_i
 | 
					 | 
				
			||||||
    redis.publish("timeline:#{timeline_id}", Oj.dump(options))
 | 
					 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def trim(type, account_id)
 | 
					  def trim(type, account_id)
 | 
				
			||||||
| 
						 | 
					@ -81,10 +76,6 @@ class FeedManager
 | 
				
			||||||
    end
 | 
					    end
 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def inline_render(target_account, template, object)
 | 
					 | 
				
			||||||
    Rabl::Renderer.new(template, object, view_path: 'app/views', format: :json, scope: InlineRablScope.new(target_account)).render
 | 
					 | 
				
			||||||
  end
 | 
					 | 
				
			||||||
 | 
					 | 
				
			||||||
  private
 | 
					  private
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def redis
 | 
					  def redis
 | 
				
			||||||
| 
						 | 
					
 | 
				
			||||||
							
								
								
									
										13
									
								
								app/lib/inline_renderer.rb
									
										
									
									
									
										Normal file
									
								
							
							
						
						
									
										13
									
								
								app/lib/inline_renderer.rb
									
										
									
									
									
										Normal file
									
								
							| 
						 | 
					@ -0,0 +1,13 @@
 | 
				
			||||||
 | 
					# frozen_string_literal: true
 | 
				
			||||||
 | 
					
 | 
				
			||||||
 | 
					class InlineRenderer
 | 
				
			||||||
 | 
					  def self.render(status, current_account, template)
 | 
				
			||||||
 | 
					    Rabl::Renderer.new(
 | 
				
			||||||
 | 
					      template,
 | 
				
			||||||
 | 
					      status,
 | 
				
			||||||
 | 
					      view_path: 'app/views',
 | 
				
			||||||
 | 
					      format: :json,
 | 
				
			||||||
 | 
					      scope: InlineRablScope.new(current_account)
 | 
				
			||||||
 | 
					    ).render
 | 
				
			||||||
 | 
					  end
 | 
				
			||||||
 | 
					end
 | 
				
			||||||
| 
						 | 
					@ -50,22 +50,22 @@ class FanOutOnWriteService < BaseService
 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def render_anonymous_payload(status)
 | 
					  def render_anonymous_payload(status)
 | 
				
			||||||
    @payload = FeedManager.instance.inline_render(nil, 'api/v1/statuses/show', status)
 | 
					    @payload = InlineRenderer.render(status, nil, 'api/v1/statuses/show')
 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def deliver_to_hashtags(status)
 | 
					  def deliver_to_hashtags(status)
 | 
				
			||||||
    Rails.logger.debug "Delivering status #{status.id} to hashtags"
 | 
					    Rails.logger.debug "Delivering status #{status.id} to hashtags"
 | 
				
			||||||
 | 
					
 | 
				
			||||||
    status.tags.pluck(:name).each do |hashtag|
 | 
					    status.tags.pluck(:name).each do |hashtag|
 | 
				
			||||||
      FeedManager.instance.broadcast("hashtag:#{hashtag}", event: 'update', payload: @payload)
 | 
					      Redis.current.publish("hashtag:#{hashtag}", Oj.dump(event: :update, payload: @payload))
 | 
				
			||||||
      FeedManager.instance.broadcast("hashtag:#{hashtag}:local", event: 'update', payload: @payload) if status.account.local?
 | 
					      Redis.current.publish("hashtag:#{hashtag}:local", Oj.dump(event: :update, payload: @payload)) if status.account.local?
 | 
				
			||||||
    end
 | 
					    end
 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def deliver_to_public(status)
 | 
					  def deliver_to_public(status)
 | 
				
			||||||
    Rails.logger.debug "Delivering status #{status.id} to public timeline"
 | 
					    Rails.logger.debug "Delivering status #{status.id} to public timeline"
 | 
				
			||||||
 | 
					
 | 
				
			||||||
    FeedManager.instance.broadcast(:public, event: 'update', payload: @payload)
 | 
					    Redis.current.publish('public', Oj.dump(event: 'update', payload: @payload))
 | 
				
			||||||
    FeedManager.instance.broadcast('public:local', event: 'update', payload: @payload) if status.account.local?
 | 
					    Redis.current.publish('public:local', Oj.dump(event: 'update', payload: @payload)) if status.account.local?
 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
end
 | 
					end
 | 
				
			||||||
| 
						 | 
					
 | 
				
			||||||
| 
						 | 
					@ -50,7 +50,7 @@ class NotifyService < BaseService
 | 
				
			||||||
  def create_notification
 | 
					  def create_notification
 | 
				
			||||||
    @notification.save!
 | 
					    @notification.save!
 | 
				
			||||||
    return unless @notification.browserable?
 | 
					    return unless @notification.browserable?
 | 
				
			||||||
    FeedManager.instance.broadcast(@recipient.id, event: 'notification', payload: FeedManager.instance.inline_render(@recipient, 'api/v1/notifications/show', @notification))
 | 
					    Redis.current.publish(@recipient.id, Oj.dump(event: :notification, payload: InlineRenderer.render(@notification, @recipient, 'api/v1/notifications/show')))
 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def send_email
 | 
					  def send_email
 | 
				
			||||||
| 
						 | 
					
 | 
				
			||||||
| 
						 | 
					@ -65,17 +65,17 @@ class RemoveStatusService < BaseService
 | 
				
			||||||
      redis.zremrangebyscore(FeedManager.instance.key(type, receiver.id), status.id, status.id)
 | 
					      redis.zremrangebyscore(FeedManager.instance.key(type, receiver.id), status.id, status.id)
 | 
				
			||||||
    end
 | 
					    end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
    FeedManager.instance.broadcast(receiver.id, event: 'delete', payload: status.id)
 | 
					    Redis.current.publish(receiver.id, Oj.dump(event: :delete, payload: status.id))
 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def remove_from_hashtags(status)
 | 
					  def remove_from_hashtags(status)
 | 
				
			||||||
    status.tags.each do |tag|
 | 
					    status.tags.each do |tag|
 | 
				
			||||||
      FeedManager.instance.broadcast("hashtag:#{tag.name}", event: 'delete', payload: status.id)
 | 
					      Redis.current.publish("hashtag:#{tag.name}", Oj.dump(event: :delete, payload: status.id))
 | 
				
			||||||
    end
 | 
					    end
 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def remove_from_public(status)
 | 
					  def remove_from_public(status)
 | 
				
			||||||
    FeedManager.instance.broadcast(:public, event: 'delete', payload: status.id)
 | 
					    Redis.current.publish('public', Oj.dump(event: :delete, payload: status.id))
 | 
				
			||||||
  end
 | 
					  end
 | 
				
			||||||
 | 
					
 | 
				
			||||||
  def redis
 | 
					  def redis
 | 
				
			||||||
| 
						 | 
					
 | 
				
			||||||
							
								
								
									
										15
									
								
								app/workers/push_update_worker.rb
									
										
									
									
									
										Normal file
									
								
							
							
						
						
									
										15
									
								
								app/workers/push_update_worker.rb
									
										
									
									
									
										Normal file
									
								
							| 
						 | 
					@ -0,0 +1,15 @@
 | 
				
			||||||
 | 
					# frozen_string_literal: true
 | 
				
			||||||
 | 
					
 | 
				
			||||||
 | 
					class PushUpdateWorker
 | 
				
			||||||
 | 
					  include Sidekiq::Worker
 | 
				
			||||||
 | 
					
 | 
				
			||||||
 | 
					  def perform(account_id, status_id)
 | 
				
			||||||
 | 
					    account = Account.find(account_id)
 | 
				
			||||||
 | 
					    status  = Status.find(status_id)
 | 
				
			||||||
 | 
					    message = InlineRenderer.render(status, account, 'api/v1/statuses/show')
 | 
				
			||||||
 | 
					
 | 
				
			||||||
 | 
					    Redis.current.publish("timeline:#{account.id}", Oj.dump(event: :update, payload: message, queued_at: (Time.now.to_f * 1000.0).to_i))
 | 
				
			||||||
 | 
					  rescue ActiveRecord::RecordNotFound
 | 
				
			||||||
 | 
					    true
 | 
				
			||||||
 | 
					  end
 | 
				
			||||||
 | 
					end
 | 
				
			||||||
		Loading…
	
	Add table
		
		Reference in a new issue