Class: Canopus::Task::Runner

Inherits:
Object
  • Object
show all
Defined in:
lib/canopus/task/runner.rb

Defined Under Namespace

Classes: Output

Constant Summary collapse

MAX_OUTPUTS =
32
MAX_CLOSERS =
32
STOP_GRACE_SECONDS =
0.25

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(scrollback:, queue_limit_bytes:, terminal_factory: nil, report: nil) ⇒ Runner

Returns a new instance of Runner.



13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
# File 'lib/canopus/task/runner.rb', line 13

def initialize(scrollback:, queue_limit_bytes:, terminal_factory: nil, report: nil)
  unless scrollback.is_a?(Integer) && scrollback.between?(0, 1_000_000)
    raise ArgumentError, "invalid task scrollback limit"
  end
  unless queue_limit_bytes.is_a?(Integer) && queue_limit_bytes.between?(65_536, 268_435_456)
    raise ArgumentError, "invalid task queue limit"
  end
  @scrollback, @queue_limit_bytes = scrollback, queue_limit_bytes
  @terminal_factory = terminal_factory || ->(**options) { Tarazed::PTY.new(**options) }
  @report = report || ->(_error) {}
  @entries, @completed = [], {}
  @close_threads = {}.compare_by_identity
  @closed_outputs = {}.compare_by_identity
  @eof_outputs = {}.compare_by_identity
  @blocked_outputs = {}.compare_by_identity
  @active_index = @sequence = @poll_index = 0
end

Instance Attribute Details

#active_index ⇒ Object (readonly)

Returns the value of attribute active_index.



11
12
13
# File 'lib/canopus/task/runner.rb', line 11

def active_index
  @active_index
end

Instance Method Details

#activate(index) ⇒ Object



68
69
70
71
72
73
74
# File 'lib/canopus/task/runner.rb', line 68

def activate(index)
  unless index.is_a?(Integer) && index.between?(0, @entries.length - 1)
    raise IndexError, "task output tab outside panel"
  end
  @active_index = index
  active
end

#active ⇒ Object



32
# File 'lib/canopus/task/runner.rb', line 32

def active = @entries[@active_index]

#close ⇒ Object

Raises:

  • (@close_error)


185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
# File 'lib/canopus/task/runner.rb', line 185

def close
  return if @closed
  @closed = true
  begin
    @close_threads.each_value(&:join)
    reap_closers
    @entries.each do |entry|
      next if @closed_outputs.key?(entry)
      begin
        close_later(entry)
      rescue StandardError => error
        record_close_error(error)
        close_now(entry)
      end
    end
  ensure
    @close_threads.each_value do |thread|
      thread.join
    rescue StandardError => error
      record_close_error(error)
    end
    @entries.each { |entry| close_now(entry) unless @closed_outputs.key?(entry) }
    @entries.clear
    @close_threads.clear
    @closed_outputs.clear
    @completed.clear
    @eof_outputs.clear
    @blocked_outputs.clear
    @sizes&.clear
  end
  raise @close_error if @close_error
  nil
end

#completed ⇒ Object



106
107
108
109
110
111
112
# File 'lib/canopus/task/runner.rb', line 106

def completed
  @entries.filter_map do |entry|
    next if @completed[entry.id] || running?(entry) || !@eof_outputs.key?(entry)
    @completed[entry.id] = true
    entry
  end
end

#drain(max_bytes:, max_seconds: 0.004) ⇒ Object

Raises:

  • (ArgumentError)


114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
# File 'lib/canopus/task/runner.rb', line 114

def drain(max_bytes:, max_seconds: 0.004)
  raise ArgumentError, "task drain limit must be positive" unless max_bytes.is_a?(Integer) && max_bytes.positive?
  unless max_seconds.is_a?(Numeric) && max_seconds.finite? && max_seconds >= 0
    raise ArgumentError, "task drain budget must be nonnegative and finite"
  end
  reap_closers
  candidates = @entries.reject { |entry| @eof_outputs.key?(entry) }
    .select { |entry| !@closed_outputs.key?(entry) || @blocked_outputs.key?(entry) }
  return false if candidates.empty?

  deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + max_seconds
  remaining, changed = max_bytes, false
  order = candidates.rotate(@poll_index % candidates.length)
  @poll_index += 1
  order.each do |entry|
    break if remaining <= 0 || Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline
    if @blocked_outputs.key?(entry)
      ready = !block_given? || yield(entry, "".b, remaining_time(deadline)) != false
      if ready
        @blocked_outputs.delete(entry)
        if @closed_outputs.key?(entry) && !@close_threads.key?(entry)
          @eof_outputs[entry] = true
          next
        end
      else
        next
      end
    end
    next if @closed_outputs.key?(entry)
    next if remaining_time(deadline).zero?
    current = entry.terminal
    method = current.method(:read)
    keywords = method.parameters.any? { |kind, _| [:key, :keyreq, :keyrest].include?(kind) }
    data = if keywords
      current.read(max_bytes: remaining,
        max_seconds: [deadline - Process.clock_gettime(Process::CLOCK_MONOTONIC), 0].max)
    else
      current.read
    end
    ready = !block_given? || yield(entry, data, remaining_time(deadline)) != false
    @blocked_outputs[entry] = true unless ready
    @eof_outputs[entry] = true if data.nil? && ready
    changed ||= !!(data && !data.empty?)
    remaining -= data.bytesize if data
  end
  changed
end

#entries ⇒ Object



31
# File 'lib/canopus/task/runner.rb', line 31

def entries = @entries.dup.freeze

#pending? ⇒ Boolean

Returns:

  • (Boolean)


162
163
164
165
166
167
168
169
170
171
172
# File 'lib/canopus/task/runner.rb', line 162

def pending?
  reap_closers
  @entries.any? do |entry|
    next running?(entry) if @eof_outputs.key?(entry)
    if @closed_outputs.key?(entry)
      next @blocked_outputs.key?(entry) || @close_threads.key?(entry)
    end
    @blocked_outputs.key?(entry) || !running?(entry) ||
      entry.terminal.respond_to?(:pending?) && entry.terminal.pending?
  end
end

#remove(index = @active_index) ⇒ Object



76
77
78
79
80
81
82
83
84
85
86
87
88
# File 'lib/canopus/task/runner.rb', line 76

def remove(index = @active_index)
  entry = @entries[index]
  return unless entry

  close_later(entry)
  @entries.delete_at(index)
  @completed.delete(entry.id)
  @eof_outputs.delete(entry)
  @blocked_outputs.delete(entry)
  @closed_outputs.delete(entry) unless @close_threads.key?(entry)
  @active_index = [index, @entries.length - 1].min.clamp(0, @entries.length)
  entry
end

#resize(columns, rows) ⇒ Object



174
175
176
177
178
179
180
181
182
183
# File 'lib/canopus/task/runner.rb', line 174

def resize(columns, rows)
  dimensions = [columns.to_i, rows.to_i]
  return if dimensions.any? { |value| value <= 0 }
  @dimensions = dimensions
  @entries.each do |entry|
    next if @closed_outputs.key?(entry) || (@sizes ||= {}.compare_by_identity)[entry.terminal] == dimensions
    entry.terminal.resize(columns: dimensions.first, rows: dimensions.last)
    @sizes[entry.terminal] = dimensions
  end
end

#run(task) ⇒ Object



35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
# File 'lib/canopus/task/runner.rb', line 35

def run(task)
  raise Error, "task runner is closed" if @closed
  validate_task!(task)
  reap_closers
  previous = @entries.find { |entry| entry.label == task.fetch("label") }
  evicted = previous || eviction_candidate
  ensure_closer_capacity! if evicted && !@closed_outputs.key?(evicted)
  dimensions = @dimensions || [100, 12]
  terminal = @terminal_factory.call(command: task.fetch("command"), cwd: task.fetch("cwd"),
    env: task.fetch("env", {}),
    columns: dimensions.first, rows: dimensions.last, scrollback: @scrollback,
    queue_limit_bytes: @queue_limit_bytes)
  @sequence += 1
  entry = Output.new(@sequence, task.fetch("label"), terminal, task.fetch("presentation")).freeze
  if evicted
    index = @entries.index { |current| current.equal?(evicted) }
    close_later(evicted)
    @completed.delete(evicted.id)
    @eof_outputs.delete(evicted)
    @blocked_outputs.delete(evicted)
    @entries[index] = entry
    @closed_outputs.delete(evicted) unless @close_threads.key?(evicted)
    @active_index = index
  else
    @entries << entry
    @active_index = @entries.length - 1
  end
  entry
rescue StandardError
  terminal&.close unless @entries.any? { |entry| entry.terminal.equal?(terminal) }
  raise
end

#running?(entry = active) ⇒ Boolean

Returns:

  • (Boolean)


98
99
100
101
102
103
104
# File 'lib/canopus/task/runner.rb', line 98

def running?(entry = active)
  return false unless entry && !@closed_outputs.key?(entry)
  terminal = entry.terminal
  !terminal.respond_to?(:alive?) || terminal.alive?
rescue IOError, SystemCallError
  false
end

#stop(entry = active) ⇒ Object



90
91
92
93
94
95
96
# File 'lib/canopus/task/runner.rb', line 90

def stop(entry = active)
  return false unless entry && @entries.any? { |current| current.equal?(entry) }
  return false unless running?(entry)

  close_later(entry, interrupt: true)
  true
end

#terminal ⇒ Object



33
# File 'lib/canopus/task/runner.rb', line 33

def terminal = active&.terminal