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_taskObject



42
43
44
# File 'lib/pg_pipeline/runtime.rb', line 42

def current_task
  Thread.current[CURRENT_TASK_KEY]
end

.deadline_for(timeout) ⇒ Object

Raises:

  • (ArgumentError)


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_nowObject



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

Returns:

  • (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

Raises:

  • (ArgumentError)


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