Class: Gienah::Transport

Inherits:
Object
  • Object
show all
Defined in:
lib/gienah/transport.rb,
sig/gienah.rbs

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(input, output, error, pid, receive) ⇒ Transport

Returns a new instance of Transport.



39
40
41
42
43
44
45
46
47
48
49
50
51
# File 'lib/gienah/transport.rb', line 39

def initialize(input, output, error, pid, receive)
  @input = input
  @output = output
  @error = error
  @pid = pid
  @receive = receive
  @write_lock = Mutex.new
  @closed = false
  @reader = Thread.new { read_loop }
  @reader.report_on_exception = false
  @stderr_reader = Thread.new { drain_stderr }
  @stderr_reader.report_on_exception = false
end

Instance Attribute Details

#pid ⇒ Integer (readonly)

Returns the value of attribute pid.

Returns:

  • (Integer)


5
6
7
# File 'lib/gienah/transport.rb', line 5

def pid
  @pid
end

Class Method Details

.open(command, cwd: nil, env: {}, policy: nil) {|arg0, arg1| ... } ⇒ Transport

Parameters:

  • command (Array[String])
  • cwd: (String, nil) (defaults to: nil)
  • env: (Hash[String, String]) (defaults to: {})
  • policy: (Object) (defaults to: nil)

Yields:

Yield Parameters:

  • arg0 (Hash[String, untyped], nil)
  • arg1 (Exception, nil)

Yield Returns:

  • (Object)

Returns:



11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
# File 'lib/gienah/transport.rb', line 11

def self.open(command, cwd: nil, env: {}, policy: nil, &receive)
  raise ArgumentError, "receiver required" unless receive
  raise ArgumentError, "command must be a nonempty Array" unless command.is_a?(Array) && !command.empty?

  input_read, input_write = IO.pipe
  output_read, output_write = IO.pipe
  error_read, error_write = IO.pipe
  [input_read, input_write, output_read, output_write, error_read, error_write].each(&:binmode)
  options = {in: input_read, out: output_write, err: error_write, close_others: true}
  options[:chdir] = cwd if cwd
  options[:new_pgroup] = true if windows?
  pid = if policy
    require "saiph"
    Saiph.spawn(command, policy: policy, **options, env: env)
  else
    Process.spawn(env, *command, **options)
  end
  [input_read, output_write, error_write].each(&:close)
  new(input_write, output_read, error_read, pid, receive)
rescue Exception
  [input_read, input_write, output_read, output_write, error_read, error_write].compact.each do |io|
    io.close unless io.closed?
  end
  Process.kill("KILL", pid) if pid
  Process.wait(pid) if pid
  raise
end

.windows? ⇒ Boolean

Returns:

  • (Boolean)


7
8
9
# File 'lib/gienah/transport.rb', line 7

def self.windows?
  /mswin|mingw/.match?(RbConfig::CONFIG.fetch("host_os"))
end

Instance Method Details

#alive? ⇒ Boolean

Returns:

  • (Boolean)


66
67
68
# File 'lib/gienah/transport.rb', line 66

def alive?
  !closed? && (@pid.nil? || process_alive?)
end

#close(grace: 0.5) ⇒ Boolean

Parameters:

  • grace: (Numeric) (defaults to: 0.5)

Returns:

  • (Boolean)


70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
# File 'lib/gienah/transport.rb', line 70

def close(grace: 0.5)
  @closed = true
  @input.close unless @input.closed?
  joined = @reader == Thread.current || @reader.join(grace)
  unless joined
    terminate("TERM")
    joined = @reader.join(grace)
  end
  unless joined
    terminate("KILL")
    @reader.join
  end
  @stderr_reader.join(grace) unless @stderr_reader == Thread.current
  [@output, @error].each { |io| io.close unless io.closed? }
  if @pid && process_alive?
    terminate("TERM")
    deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + grace
    until Process.waitpid(@pid, Process::WNOHANG)
      break if Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline

      sleep(0.01)
    end
    if process_alive?
      terminate("KILL")
      Process.wait(@pid)
    end
  end
  true
rescue IOError, Errno::ECHILD
  true
end

#write(message, max_size: Protocol::MAX_MESSAGE) ⇒ nil

Parameters:

  • message (Hash[String, untyped])
  • max_size: (Integer) (defaults to: Protocol::MAX_MESSAGE)

Returns:

  • (nil)


53
54
55
56
57
58
59
60
61
62
63
64
# File 'lib/gienah/transport.rb', line 53

def write(message, max_size: Protocol::MAX_MESSAGE)
  frame = Protocol.frame(message, max_size: max_size)
  @write_lock.synchronize do
    raise Error, "transport is closed" if closed?

    @input.write(frame)
    @input.flush
  end
  nil
rescue IOError, Errno::EPIPE => error
  raise Error, "transport write failed: #{error.message}"
end