Class: Gritz::Transport::Native
- Inherits:
-
Object
- Object
- Gritz::Transport::Native
- 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
-
#bind(listener_spec = @config.bind) ⇒ Integer
Creates C-core resources only when binding, so construction is safe before fork.
- #drain! ⇒ Object
-
#initialize(config:, dispatcher:, logger:) ⇒ Native
constructor
A new instance of Native.
- #kill ⇒ Object
- #running? ⇒ Boolean
- #start ⇒ Object
- #stats ⇒ Object
- #stop(deadline:) ⇒ Object
- #update_health(ready:, checks: {}) ⇒ Object
- #wait ⇒ Object
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.}. 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.
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
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 |