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, Notification, Queue, Semaphore, Task, TimeoutError
Constant Summary
collapse
- CURRENT_TASK_KEY =
:pg_pipeline_current_task
- WAITER_KEY =
:pg_pipeline_waiter
Class Method Summary
collapse
Class Method Details
.build_waiter(blocker) ⇒ Object
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
|
# File 'lib/pg_pipeline/runtime.rb', line 29
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
25
26
27
|
# File 'lib/pg_pipeline/runtime.rb', line 25
def current_task
Thread.current[CURRENT_TASK_KEY]
end
|
.park(blocker, waiter) ⇒ Object
47
48
49
50
51
52
53
54
55
56
57
58
59
60
|
# File 'lib/pg_pipeline/runtime.rb', line 47
def park(blocker, waiter)
scheduler = waiter[:scheduler]
task = current_task
task&.enter_block(waiter)
until waiter[:ready] || yield
scheduler.block(blocker, nil)
task&.raise_if_cancelled!
end
nil
ensure
waiter[:blocker] = nil
task&.exit_block
end
|
.scheduler! ⇒ Object
21
22
23
|
# File 'lib/pg_pipeline/runtime.rb', line 21
def scheduler!
Fiber.scheduler or raise Error, "operation requires an active Fiber scheduler"
end
|
.spawn(name: nil, &block) ⇒ Object
17
18
19
|
# File 'lib/pg_pipeline/runtime.rb', line 17
def spawn(name: nil, &block)
Task.spawn(name: name, &block)
end
|
.wake(waiter, blocker) ⇒ Object
85
86
87
88
89
90
91
|
# File 'lib/pg_pipeline/runtime.rb', line 85
def wake(waiter, blocker)
waiter[:ready] = true
fiber = waiter[:fiber]
return unless fiber.alive?
waiter[:scheduler].unblock(blocker, fiber)
end
|
.wake_dequeued(waiter, blocker) ⇒ Object
71
72
73
74
|
# File 'lib/pg_pipeline/runtime.rb', line 71
def wake_dequeued(waiter, blocker)
waiter[:queued] = false
wake(waiter, blocker)
end
|
.with_timeout(duration) ⇒ Object
76
77
78
79
80
81
82
83
|
# File 'lib/pg_pipeline/runtime.rb', line 76
def with_timeout(duration)
return yield if duration.nil?
timeout = Float(duration)
raise ArgumentError, "timeout must be non-negative and finite" unless timeout.finite? && timeout >= 0
::Timeout.timeout(timeout, TimeoutError) { yield }
end
|
.with_waiter(blocker, waiters) ⇒ Object
62
63
64
65
66
67
68
69
|
# File 'lib/pg_pipeline/runtime.rb', line 62
def with_waiter(blocker, waiters)
waiter = build_waiter(blocker)
waiter[:queued] = true
waiters << waiter
yield waiter
ensure
waiters.delete(waiter) if waiter[:queued]
end
|