Class: Gritz::Transport::Native

Inherits:
Object
  • Object
show all
Defined in:
/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb,
/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native/bridge.rb,
/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native/call.rb,
/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native/cancellation.rb,
/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native/client.rb,
/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native/client_invocation.rb,
/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native/health.rb,
/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native/reflection.rb,
/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native/server.rb

Overview

Runs generated services through Gritz's transport-independent dispatcher.

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(config:, dispatcher:, logger:) ⇒ Native

Returns a new instance of Native.



22
23
24
25
26
27
28
29
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 22

def initialize(config:, dispatcher:, logger:)
  @config = config
  @dispatcher = dispatcher
  @logger = logger
  @lock = Mutex.new
  @inflight = {}.compare_by_identity
  @requests_total = @rejected_total = 0
end

Class Method Details

.capabilities ⇒ Object



11
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 11

def self.capabilities = Set[:unary, :client_streaming, :server_streaming, :bidi, :reuseport, :health, :tls, :mtls, :reflection].freeze

.postfork_child ⇒ Object



20
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 20

def self.postfork_child = GRPC.postfork_child

.postfork_parent ⇒ Object



19
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 19

def self.postfork_parent = GRPC.postfork_parent

.prefork ⇒ Object



13
14
15
16
17
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 13

def self.prefork
  GRPC.prefork
rescue RuntimeError => e
  raise "#{e.message}. fork_mode :grpc_fork_support requires GRPC_ENABLE_FORK_SUPPORT=1 before requiring grpc."
end

Instance Method Details

#bind(listener_spec = @config.bind) ⇒ Integer

Creates C-core resources only when binding, so construction is safe before fork.

Returns:

  • (Integer) —

    the port selected by C-core



33
34
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
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 33

def bind(listener_spec = @config.bind)
  raise ArgumentError, "transport is already bound" if @server
  raise ArgumentError, "at least one controller must be registered" if @dispatcher.router.routes.empty?

  credentials = server_credentials
  created_server = @server = Server.new(
    pool_size: @config.threads,
    # grpc 1.83 accepts but ignores this value: busy pools reject immediately.
    max_waiting_requests: @config.max_waiting_requests,
    poll_period: @config.shutdown_timeout,
    pool_keep_alive: 0,
    server_args: @config.server_args,
    on_rejected: -> { @lock.synchronize { @rejected_total += 1 } }
  )
  @dispatcher.router.routes.values.group_by(&:service_class).each do |service_class, descriptors|
    @server.handle(Bridge.build(service_class, descriptors, self))
  end
  @health = Health.new(@dispatcher.router.routes.values.map(&:service).uniq)
  @server.handle(@health)
  if @config.reflection
    require_relative "native/reflection"
    services = @dispatcher.router.routes.values.map(&:service).uniq + [Health.service_name]
    Reflection.build(services).each { |service| @server.handle(service) }
  end
  @port = @server.add_http2_port(listener_spec, credentials)
  raise ArgumentError, "could not bind #{listener_spec}" unless @port.positive?

  @port
rescue StandardError
  created_server&.close_unstarted
  @server = nil if created_server
  raise
end

#drain! ⇒ Object



108
109
110
111
112
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 108

def drain!
  @draining = true
  @health&.drain!
  self
end

#kill ⇒ Object



100
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 100

def kill = stop(deadline: Time.now)

#running? ⇒ Boolean

Returns:

  • (Boolean)


86
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 86

def running? = @thread&.alive? && @server.running?

#start ⇒ Object



67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 67

def start
  raise ArgumentError, "bind must be called before start" unless @server
  raise ArgumentError, "transport is already started" if @thread

  @thread = Thread.new { @server.run }
  if @server.wait_till_running(5)
    refresh_health
    return self
  end

  @thread.value unless @thread.alive?
  kill
  raise "gRPC server did not start within 5 seconds"
rescue StandardError
  @server&.close_unstarted
  raise
end

#stats ⇒ Object



125
126
127
128
129
130
131
132
133
134
135
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 125

def stats
  busy = @server&.busy_threads || 0
  @lock.synchronize do
    oldest = @inflight.values.min
    {
      inflight: @inflight.size, busy: busy, capacity: @config.threads,
      rejected_total: @rejected_total, requests_total: @requests_total,
      oldest_inflight_age: oldest ? Process.clock_gettime(Process::CLOCK_MONOTONIC) - oldest : 0
    }
  end
end

#stop(deadline:) ⇒ Object



88
89
90
91
92
93
94
95
96
97
98
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 88

def stop(deadline:)
  drain!
  unless @thread
    @server&.close_unstarted
    @server = nil
    return
  end

  @server.stop(deadline:)
  wait
end

#update_health(ready:, checks: {}) ⇒ Object



102
103
104
105
106
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 102

def update_health(ready:, checks: {})
  healthy = ready && checks.values.all? && !@draining
  @health&.update(healthy)
  healthy
end

#wait ⇒ Object



85
# File '/home/runner/work/gritz/gritz/sources/gritz-native/lib/gritz/transport/native.rb', line 85

def wait = @thread&.value