@@ -24,6 +24,7 @@ module Mongo
2424 # @since 2.0.0
2525 class Cluster
2626 extend Forwardable
27+ include Monitoring ::Publishable
2728 include Event ::Subscriber
2829 include Loggable
2930
@@ -45,6 +46,9 @@ class Cluster
4546 # @return [ Hash ] The options hash.
4647 attr_reader :options
4748
49+ # @return [ Monitoring ] monitoring The monitoring.
50+ attr_reader :monitoring
51+
4852 # @return [ Object ] The cluster topology.
4953 attr_reader :topology
5054
@@ -54,7 +58,8 @@ class Cluster
5458 # @since 2.4.0
5559 attr_reader :app_metadata
5660
57- def_delegators :topology , :replica_set? , :replica_set_name , :sharded? , :single? , :unknown?
61+ def_delegators :topology , :replica_set? , :replica_set_name , :sharded? ,
62+ :single? , :unknown? , :member_discovered
5863 def_delegators :@cursor_reaper , :register_cursor , :schedule_kill_cursor , :unregister_cursor
5964
6065 # Determine if this cluster of servers is equal to another object. Checks the
@@ -89,7 +94,6 @@ def add(host)
8994 address = Address . new ( host )
9095 if !addresses . include? ( address )
9196 if addition_allowed? ( address )
92- log_debug ( "Adding #{ address . to_s } to the cluster." )
9397 @update_lock . synchronize { @addresses . push ( address ) }
9498 server = Server . new ( address , self , @monitoring , event_listeners , options )
9599 @update_lock . synchronize { @servers . push ( server ) }
@@ -98,6 +102,34 @@ def add(host)
98102 end
99103 end
100104
105+ # Determine if the cluster would select a readable server for the
106+ # provided read preference.
107+ #
108+ # @example Is a readable server present?
109+ # topology.has_readable_server?(server_selector)
110+ #
111+ # @param [ ServerSelector ] server_selector The server
112+ # selector.
113+ #
114+ # @return [ true, false ] If a readable server is present.
115+ #
116+ # @since 2.4.0
117+ def has_readable_server? ( server_selector = nil )
118+ topology . has_readable_server? ( self , server_selector )
119+ end
120+
121+ # Determine if the cluster would select a writable server.
122+ #
123+ # @example Is a writable server present?
124+ # topology.has_writable_server?
125+ #
126+ # @return [ true, false ] If a writable server is present.
127+ #
128+ # @since 2.4.0
129+ def has_writable_server?
130+ topology . has_writable_server? ( self )
131+ end
132+
101133 # Instantiate the new cluster.
102134 #
103135 # @api private
@@ -119,16 +151,26 @@ def initialize(seeds, monitoring, options = Options::Redacted.new)
119151 @event_listeners = Event ::Listeners . new
120152 @options = options . freeze
121153 @app_metadata ||= AppMetadata . new ( self )
122- @topology = Topology . initial ( seeds , options )
123154 @update_lock = Mutex . new
124155 @pool_lock = Mutex . new
156+ @topology = Topology . initial ( seeds , monitoring , options )
157+
158+ publish_sdam_event (
159+ Monitoring ::TOPOLOGY_OPENING ,
160+ Monitoring ::Event ::TopologyOpening . new ( @topology )
161+ )
125162
126163 subscribe_to ( Event ::STANDALONE_DISCOVERED , Event ::StandaloneDiscovered . new ( self ) )
127164 subscribe_to ( Event ::DESCRIPTION_CHANGED , Event ::DescriptionChanged . new ( self ) )
128- subscribe_to ( Event ::PRIMARY_ELECTED , Event ::PrimaryElected . new ( self ) )
165+ subscribe_to ( Event ::MEMBER_DISCOVERED , Event ::MemberDiscovered . new ( self ) )
129166
130167 seeds . each { |seed | add ( seed ) }
131168
169+ publish_sdam_event (
170+ Monitoring ::TOPOLOGY_CHANGED ,
171+ Monitoring ::Event ::TopologyChanged . new ( @topology , @topology )
172+ ) if @servers . size > 1
173+
132174 @cursor_reaper = CursorReaper . new
133175 @cursor_reaper . run!
134176
@@ -264,11 +306,14 @@ def standalone_discovered
264306 #
265307 # @since 2.0.0
266308 def remove ( host )
267- log_debug ( "#{ host } being removed from the cluster." )
268309 address = Address . new ( host )
269310 removed_servers = @servers . select { |s | s . address == address }
270311 @update_lock . synchronize { @servers = @servers - removed_servers }
271312 removed_servers . each { |server | server . disconnect! } if removed_servers
313+ publish_sdam_event (
314+ Monitoring ::SERVER_CLOSED ,
315+ Monitoring ::Event ::ServerClosed . new ( address , topology )
316+ )
272317 @update_lock . synchronize { @addresses . reject! { |addr | addr == address } }
273318 end
274319
0 commit comments