Skip to content

Commit cd55d8b

Browse files
added:
* worker to check fiber timeouts * timeout_after implementation for fiber scheduler
1 parent 86685d8 commit cd55d8b

1 file changed

Lines changed: 78 additions & 0 deletions

File tree

lib/rage/fiber_scheduler.rb

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,16 +4,29 @@
44

55
class Rage::FiberScheduler
66
MAX_READ = 65536
7+
TIMEOUT_WORKER_INTERVAL = 100 # miliseconds
78

89
def initialize
910
@root_fiber = Fiber.current
1011
@dns_cache = {}
12+
13+
@alive_fibers = {}
14+
@timeout_mutex = Mutex.new
15+
16+
start_timeout_worker
1117
end
1218

1319
def io_wait(io, events, timeout = nil)
1420
f = Fiber.current
1521
::Iodine::Scheduler.attach(io.fileno, events, timeout&.ceil) { |err| f.resume(err) }
1622

23+
timeout_deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout
24+
@alive_fibers[f.__get_id] = {
25+
fiber: f,
26+
timeout_deadline: timeout_deadline,
27+
exception_class: RageTimeout,
28+
}
29+
1730
err = Fiber.defer(io.fileno)
1831
if err == false || (err && err < 0)
1932
err
@@ -79,6 +92,30 @@ def kernel_sleep(duration = nil)
7992

8093
# result
8194
# end
95+
def timeout_after(duration, exception_class = Timeout::Error, *exception_arguments, &block)
96+
fiber = Fiber.current
97+
timeout_deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + duration
98+
99+
p "duration #{duration}"
100+
p "fiber id #{fiber.__get_id}"
101+
102+
@timeout_mutex.synchronize do
103+
@alive_fibers[fiber.__get_id] = {
104+
fiber: fiber,
105+
timeout_deadline: timeout_deadline,
106+
exception_class: exception_class,
107+
exception_arguments: exception_arguments,
108+
}
109+
end
110+
111+
begin
112+
block.call
113+
ensure
114+
@timeout_mutex.synchronize do
115+
@alive_fibers.delete(fiber.__get_id)
116+
end
117+
end
118+
end
82119

83120
def address_resolve(hostname)
84121
@dns_cache[hostname] ||= begin
@@ -145,4 +182,45 @@ def fiber(&block)
145182
def close
146183
::Iodine::Scheduler.close
147184
end
185+
186+
private
187+
188+
def start_timeout_worker
189+
return unless ::Iodine.running?
190+
191+
::Iodine.run_every(Rage::FiberScheduler::TIMEOUT_WORKER_INTERVAL) do
192+
@timeout_mutex.synchronize do
193+
check_timeouts
194+
end
195+
end
196+
end
197+
198+
def check_timeouts
199+
p @alive_fibers.count
200+
201+
@alive_fibers.delete_if do |_, fiber_hash|
202+
current_time = Process.clock_gettime(Process::CLOCK_MONOTONIC)
203+
204+
p current_time
205+
p "deadline #{fiber_hash[:timeout_deadline]}"
206+
p "fiber id #{fiber_hash[:fiber].__get_id}"
207+
208+
return false if current_time < fiber_hash[:timeout_deadline]
209+
210+
p 'after'
211+
212+
fiber = fiber_hash[:fiber]
213+
# unblock(nil, fiber)
214+
215+
# if fiber.alive?
216+
fiber.raise(RageTimeout)
217+
# else
218+
# fiber.kill
219+
# end
220+
221+
true
222+
end
223+
end
148224
end
225+
226+
class RageTimeout < StandardError; end

0 commit comments

Comments
 (0)