Module: PgPipeline::Runtime
- Defined in:
- lib/pg_pipeline/runtime.rb,
lib/pg_pipeline/runtime/task.rb,
lib/pg_pipeline/runtime/queue.rb,
lib/pg_pipeline/runtime/semaphore.rb,
lib/pg_pipeline/runtime/notification.rb
Defined Under Namespace
Classes: Cancel, Deadline, Notification, Queue, Semaphore, Task, TimeoutError
Constant Summary
collapse
- CURRENT_TASK_KEY =
:pg_pipeline_current_task
- WAITER_KEY =
:pg_pipeline_waiter
- DEADLINE_MESSAGE =
"execution expired"
- MISSING_TIMEOUT_HOOK_WARNING =
<<~MESSAGE
pg_pipeline: %s does not implement #timeout_after.
Falling back to stdlib Timeout, which uses Thread#raise and can deliver
the timeout to an unrelated fiber. Use a scheduler that implements
#timeout_after (async, itsi-scheduler) or avoid Runtime.with_timeout
on this host.
MESSAGE
Class Method Summary
collapse
Class Method Details
.arm_deadline(seconds, &block) ⇒ Object
137
138
139
140
141
142
143
144
145
|
# File 'lib/pg_pipeline/runtime.rb', line 137
def arm_deadline(seconds, &block)
target = Fiber.scheduler
if native_timeouts?(target)
return target.timeout_after(seconds, Deadline, DEADLINE_MESSAGE, &block)
end
warn_missing_timeout_hook(target) if target
::Timeout.timeout(seconds, Deadline, DEADLINE_MESSAGE, &block)
end
|
.build_waiter(blocker) ⇒ Object
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
|
# File 'lib/pg_pipeline/runtime.rb', line 46
def build_waiter(blocker)
waiter = Thread.current[WAITER_KEY] ||= {
fiber: nil, scheduler: nil, ready: false, blocker: nil
}
if waiter[:blocker]
raise Error, "waiter already parked on #{waiter[:blocker].class}; " \
"a fiber may only park in one place at a time"
end
waiter[:fiber] = Fiber.current
waiter[:scheduler] = scheduler!
waiter[:ready] = false
waiter[:queued] = false
waiter[:blocker] = blocker
waiter
end
|
.current_task ⇒ Object
42
43
44
|
# File 'lib/pg_pipeline/runtime.rb', line 42
def current_task
Thread.current[CURRENT_TASK_KEY]
end
|
.deadline_for(timeout) ⇒ Object
115
116
117
118
119
120
121
122
|
# File 'lib/pg_pipeline/runtime.rb', line 115
def deadline_for(timeout)
return nil if timeout.nil?
seconds = Float(timeout)
raise ArgumentError, "timeout must be non-negative and finite" unless seconds.finite? && seconds >= 0
monotonic_now + seconds
end
|
.monotonic_now ⇒ Object
111
112
113
|
# File 'lib/pg_pipeline/runtime.rb', line 111
def monotonic_now
Process.clock_gettime(Process::CLOCK_MONOTONIC)
end
|
.native_timeouts?(target = Fiber.scheduler) ⇒ Boolean
38
39
40
|
# File 'lib/pg_pipeline/runtime.rb', line 38
def native_timeouts?(target = Fiber.scheduler)
!target.nil? && target.respond_to?(:timeout_after)
end
|
.park(blocker, waiter, deadline = nil) ⇒ Object
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
|
# File 'lib/pg_pipeline/runtime.rb', line 64
def park(blocker, waiter, deadline = nil)
scheduler = waiter[:scheduler]
task = current_task
task&.enter_block(waiter)
until waiter[:ready] || yield
remaining = deadline ? (deadline - monotonic_now) : nil
return false if deadline && remaining <= 0
scheduler.block(blocker, remaining)
task&.raise_if_cancelled!
end
task&.raise_if_cancelled!
true
ensure
waiter[:blocker] = nil
task&.exit_block
end
|
.scheduler! ⇒ Object
34
35
36
|
# File 'lib/pg_pipeline/runtime.rb', line 34
def scheduler!
Fiber.scheduler or raise Error, "operation requires an active Fiber scheduler"
end
|
.spawn(name: nil, &block) ⇒ Object
30
31
32
|
# File 'lib/pg_pipeline/runtime.rb', line 30
def spawn(name: nil, &block)
Task.spawn(name: name, &block)
end
|
.wake(waiter, blocker) ⇒ Object
100
101
102
103
104
105
106
107
108
109
|
# File 'lib/pg_pipeline/runtime.rb', line 100
def wake(waiter, blocker)
return false if waiter[:ready]
waiter[:ready] = true
fiber = waiter[:fiber]
return false unless fiber.alive?
waiter[:scheduler].unblock(blocker, fiber)
true
end
|
.wake_dequeued(waiter, blocker) ⇒ Object
95
96
97
98
|
# File 'lib/pg_pipeline/runtime.rb', line 95
def wake_dequeued(waiter, blocker)
waiter[:queued] = false
wake(waiter, blocker)
end
|
.warn_missing_timeout_hook(target) ⇒ Object
147
148
149
150
151
152
153
|
# File 'lib/pg_pipeline/runtime.rb', line 147
def warn_missing_timeout_hook(target)
key = target.class
return if @warned_schedulers[key]
@warned_schedulers[key] = true
warn(format(MISSING_TIMEOUT_HOOK_WARNING, key))
end
|
.with_timeout(duration) ⇒ Object
124
125
126
127
128
129
130
131
132
133
134
135
|
# File 'lib/pg_pipeline/runtime.rb', line 124
def with_timeout(duration)
return yield if duration.nil?
seconds = Float(duration)
raise ArgumentError, "timeout must be non-negative and finite" unless seconds.finite? && seconds >= 0
begin
arm_deadline(seconds) { yield }
rescue Deadline => e
raise TimeoutError, e.message
end
end
|
.with_waiter(blocker, waiters) ⇒ Object
83
84
85
86
87
88
89
90
91
92
93
|
# File 'lib/pg_pipeline/runtime.rb', line 83
def with_waiter(blocker, waiters)
waiter = build_waiter(blocker)
waiter[:queued] = true
waiters << waiter
yield waiter
ensure
if waiter[:queued]
waiters.delete(waiter)
waiter[:queued] = false
end
end
|