Module: Einhorn::Event

Defined in:
lib/einhorn/event/loop_breaker.rb,
lib/einhorn/event.rb,
lib/einhorn/event/timer.rb,
lib/einhorn/event/ack_timer.rb,
lib/einhorn/event/connection.rb,
lib/einhorn/event/persistent.rb,
lib/einhorn/event/command_server.rb,
lib/einhorn/event/abstract_text_descriptor.rb

Overview

TODO: set lots of cloexecs

Defined Under Namespace

Modules: Persistent Classes: ACKTimer, AbstractTextDescriptor, CommandServer, Connection, LoopBreaker, Timer

Constant Summary collapse

@@loopbreak_reader =
nil
@@loopbreak_writer =
nil
@@default_timeout =
nil
@@signal_actions =
[]
@@readable =
{}
@@writeable =
{}
@@connections =
{}
@@timers =
{}

Class Method Summary collapse

Class Method Details

.break_loopObject



161
162
163
164
165
166
167
168
# File 'lib/einhorn/event.rb', line 161

def self.break_loop
  Einhorn.log_debug("Breaking the loop")
  begin
    @@loopbreak_writer.write_nonblock("a")
  rescue Errno::EWOULDBLOCK, Errno::EAGAIN
    Einhorn.log_error("Loop break pipe is full -- probably means that we are quite backlogged")
  end
end

.close_allObject



24
25
26
27
28
29
30
31
32
# File 'lib/einhorn/event.rb', line 24

def self.close_all
  @@loopbreak_reader.close
  @@loopbreak_writer.close
  (@@readable.values + @@writeable.values).each do |descriptors|
    descriptors.each do |descriptor|
      descriptor.close
    end
  end
end

.close_all_for_workerObject



34
35
36
# File 'lib/einhorn/event.rb', line 34

def self.close_all_for_worker
  close_all
end

.connectionsObject



98
99
100
# File 'lib/einhorn/event.rb', line 98

def self.connections
  @@connections.values
end

.default_timeoutObject



174
175
176
# File 'lib/einhorn/event.rb', line 174

def self.default_timeout
  @@default_timeout
end

.default_timeout=(val) ⇒ Object



170
171
172
# File 'lib/einhorn/event.rb', line 170

def self.default_timeout=(val)
  @@default_timeout = (val.to_i == 0) ? nil : val.to_i
end

.deregister_connection(fd) ⇒ Object



94
95
96
# File 'lib/einhorn/event.rb', line 94

def self.deregister_connection(fd)
  @@connections.delete(fd)
end

.deregister_readable(reader) ⇒ Object



59
60
61
62
63
# File 'lib/einhorn/event.rb', line 59

def self.deregister_readable(reader)
  readers = @@readable[reader.to_io]
  readers.delete(reader)
  @@readable.delete(reader.to_io) if readers.length == 0
end

.deregister_timer(timer) ⇒ Object



107
108
109
110
111
# File 'lib/einhorn/event.rb', line 107

def self.deregister_timer(timer)
  timers = @@timers[timer.expires_at]
  timers.delete(timer)
  @@timers.delete(timer.expires_at) if timers.length == 0
end

.deregister_writeable(writer) ⇒ Object



76
77
78
79
80
# File 'lib/einhorn/event.rb', line 76

def self.deregister_writeable(writer)
  writers = @@writeable[writer.to_io]
  writers.delete(writer)
  @@writeable.delete(writer.to_io) if writers.length == 0
end

.initObject



14
15
16
17
18
19
20
21
22
# File 'lib/einhorn/event.rb', line 14

def self.init
  readable, writeable = Einhorn::Compat.pipe

  @@loopbreak_reader = LoopBreaker.open(readable)
  @@loopbreak_writer = writeable

  Einhorn::Compat.cloexec!(readable, true)
  Einhorn::Compat.cloexec!(writeable, true)
end

.loop_onceObject



113
114
115
116
117
# File 'lib/einhorn/event.rb', line 113

def self.loop_once
  run_signal_actions
  run_selectables
  run_timers
end

.persistent_descriptorsObject



38
39
40
41
42
# File 'lib/einhorn/event.rb', line 38

def self.persistent_descriptors
  descriptor_sets = @@readable.values + @@writeable.values + @@timers.values
  descriptors = descriptor_sets.inject { |a, b| a | b }
  descriptors.select { |descriptor| Einhorn::Event::Persistent.persistent?(descriptor) }
end

.readable_fdsObject



65
66
67
68
69
# File 'lib/einhorn/event.rb', line 65

def self.readable_fds
  readers = @@readable.keys
  Einhorn.log_debug("Readable fds are #{readers.inspect}")
  readers
end

.register_connection(connection, fd) ⇒ Object



90
91
92
# File 'lib/einhorn/event.rb', line 90

def self.register_connection(connection, fd)
  @@connections[fd] = connection
end

.register_readable(reader) ⇒ Object



54
55
56
57
# File 'lib/einhorn/event.rb', line 54

def self.register_readable(reader)
  @@readable[reader.to_io] ||= Set.new
  @@readable[reader.to_io] << reader
end

.register_signal_action(&blk) ⇒ Object



50
51
52
# File 'lib/einhorn/event.rb', line 50

def self.register_signal_action(&blk)
  @@signal_actions << blk
end

.register_timer(timer) ⇒ Object



102
103
104
105
# File 'lib/einhorn/event.rb', line 102

def self.register_timer(timer)
  @@timers[timer.expires_at] ||= Set.new
  @@timers[timer.expires_at] << timer
end

.register_writeable(writer) ⇒ Object



71
72
73
74
# File 'lib/einhorn/event.rb', line 71

def self.register_writeable(writer)
  @@writeable[writer.to_io] ||= Set.new
  @@writeable[writer.to_io] << writer
end

.restore_persistent_descriptors(persistent_descriptors) ⇒ Object



44
45
46
47
48
# File 'lib/einhorn/event.rb', line 44

def self.restore_persistent_descriptors(persistent_descriptors)
  persistent_descriptors.each do |descriptor_state|
    Einhorn::Event::Persistent.from_state(descriptor_state)
  end
end

.run_selectablesObject



138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
# File 'lib/einhorn/event.rb', line 138

def self.run_selectables
  time = timeout
  Einhorn.log_debug("Loop timeout is #{time.inspect}")
  # Time's already up
  return if time && time < 0

  readable, writeable, _ = IO.select(readable_fds, writeable_fds, nil, time)
  (readable || []).each do |io|
    @@readable[io].each { |reader| reader.notify_readable }
  end

  (writeable || []).each do |io|
    @@writeable[io].each { |writer| writer.notify_writeable }
  end
end

.run_signal_actionsObject



128
129
130
131
132
133
134
135
136
# File 'lib/einhorn/event.rb', line 128

def self.run_signal_actions
  # Note thah @@signal_actions can be mutated in the signal
  # handlers. Since it's just an array we push to/shift from, we
  # can be sure there's no race (such as adding hash keys during
  # iteration.)
  while (blk = @@signal_actions.shift)
    blk.call
  end
end

.run_timersObject



154
155
156
157
158
159
# File 'lib/einhorn/event.rb', line 154

def self.run_timers
  @@timers.select { |expires_at, _| expires_at <= Time.now }.each do |expires_at, timers|
    # Going to be modifying the set, so let's dup it.
    timers.dup.each { |timer| timer.ring! }
  end
end

.timeoutObject



119
120
121
122
123
124
125
126
# File 'lib/einhorn/event.rb', line 119

def self.timeout
  # (expires_at of the next timer) - now
  if (expires_at = @@timers.keys.min)
    expires_at - Time.now
  else
    @@default_timeout
  end
end

.writeable_fdsObject



82
83
84
85
86
87
88
# File 'lib/einhorn/event.rb', line 82

def self.writeable_fds
  writers = @@writeable.select do |io, writers|
    writers.any? { |writer| writer.write_pending? }
  end.map { |io, writers| io }
  Einhorn.log_debug("Writeable fds are #{writers.inspect}")
  writers
end