Class: SolidRedis::Subscription

Inherits:
Client
  • Object
show all
Defined in:
lib/solid_redis/subscription.rb

Overview

A dedicated Pub/Sub connection.

Once SUBSCRIBE has been sent, a Redis connection stops answering regular commands and turns into a stream of push messages. A Subscription therefore owns its own socket, is never taken from a pool, and belongs to the Ractor that created it. Channel and pattern lists are tracked locally so that the connection can be re-established and re-subscribed after a network error.

Defined Under Namespace

Classes: Message

Constant Summary collapse

SUBSCRIBE_COMMANDS =
Ractor.make_shareable({
  channels: %w[SUBSCRIBE UNSUBSCRIBE],
  patterns: %w[PSUBSCRIBE PUNSUBSCRIBE],
  shards: %w[SSUBSCRIBE SUNSUBSCRIBE],
})

Instance Attribute Summary

Attributes inherited from Client

#config

Instance Method Summary collapse

Methods inherited from Client

#close, #connected?, #server_url

Constructor Details

#initialize(config, name: nil) ⇒ Subscription

Returns a new instance of Subscription.



24
25
26
27
28
29
# File 'lib/solid_redis/subscription.rb', line 24

def initialize(config, name: nil)
  super
  @channels = []
  @patterns = []
  @shards = []
end

Instance Method Details

#call ⇒ Object Also known as: call_v, pipelined, blocking_call, blocking_call_v

Raises:



94
95
96
# File 'lib/solid_redis/subscription.rb', line 94

def call(*)
  raise Error, "Regular commands are not available on a Pub/Sub connection"
end

#each_message(timeout: nil) ⇒ Object

Yields every incoming Message. When timeout is given, yields nil each time it elapses so the caller can check a stop condition.



88
89
90
91
92
# File 'lib/solid_redis/subscription.rb', line 88

def each_message(timeout: nil)
  return enum_for(__method__, timeout: timeout) unless block_given?

  loop { yield next_message(timeout: timeout) }
end

#next_message(timeout: nil) ⇒ Object

Returns the next Message, or nil when timeout seconds elapse without one. A nil timeout waits forever. Connection errors trigger a reconnect and a re-subscription according to reconnect_attempts.



71
72
73
74
75
76
77
78
79
80
81
82
83
84
# File 'lib/solid_redis/subscription.rb', line 71

def next_message(timeout: nil)
  attempts = 0
  loop do
    ensure_connected
    return unless @reader.wait_readable(timeout)

    return decode(@reader.read)
  rescue ConnectionError, IO::WaitReadable, IO::WaitWritable, SystemCallError => error
    handle_error(error)
    raise error if attempts >= config.reconnect_attempts

    attempts += 1
  end
end

#ping(payload = nil) ⇒ Object



55
56
57
58
# File 'lib/solid_redis/subscription.rb', line 55

def ping(payload = nil)
  transmit { payload ? ["PING", payload] : ["PING"] }
  self
end

#psubscribe(*patterns) ⇒ Object



35
36
37
# File 'lib/solid_redis/subscription.rb', line 35

def psubscribe(*patterns)
  change(:patterns, 0, patterns)
end

#punsubscribe(*patterns) ⇒ Object



47
48
49
# File 'lib/solid_redis/subscription.rb', line 47

def punsubscribe(*patterns)
  change(:patterns, 1, patterns)
end

#ssubscribe(*channels) ⇒ Object



39
40
41
# File 'lib/solid_redis/subscription.rb', line 39

def ssubscribe(*channels)
  change(:shards, 0, channels)
end

#subscribe(*channels) ⇒ Object



31
32
33
# File 'lib/solid_redis/subscription.rb', line 31

def subscribe(*channels)
  change(:channels, 0, channels)
end

#subscribed? ⇒ Boolean

Returns:

  • (Boolean)


64
65
66
# File 'lib/solid_redis/subscription.rb', line 64

def subscribed?
  !(@channels.empty? && @patterns.empty? && @shards.empty?)
end

#subscriptions ⇒ Object



60
61
62
# File 'lib/solid_redis/subscription.rb', line 60

def subscriptions
  { channels: @channels.dup, patterns: @patterns.dup, shards: @shards.dup }
end

#sunsubscribe(*channels) ⇒ Object



51
52
53
# File 'lib/solid_redis/subscription.rb', line 51

def sunsubscribe(*channels)
  change(:shards, 1, channels)
end

#unsubscribe(*channels) ⇒ Object



43
44
45
# File 'lib/solid_redis/subscription.rb', line 43

def unsubscribe(*channels)
  change(:channels, 1, channels)
end