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_taskObject



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

Raises:

  • (ArgumentError)


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