Class: Hatchet::WorkerRuntime::ActionListener
- Inherits:
-
Object
- Object
- Hatchet::WorkerRuntime::ActionListener
- Defined in:
- lib/hatchet/worker/action_listener.rb,
sig/hatchet/worker/action_listener.rbs
Overview
Listens for action assignments from the Hatchet dispatcher via gRPC streaming.
The action listener establishes a server-streaming gRPC connection (ListenV2), receives action assignments (task executions), and forwards them to the runner for execution.
Includes retry/reconnect logic with exponential backoff and a heartbeat thread.
Constant Summary collapse
- MAX_RETRIES =
15- BASE_BACKOFF_SECONDS =
1- MAX_BACKOFF_SECONDS =
30- HEALTHY_CONNECTION_THRESHOLD =
seconds
30- HEARTBEAT_INTERVAL =
seconds
4- MAX_MISSED_HEARTBEATS =
3
Instance Method Summary collapse
-
#initialize(dispatcher_client:, worker_id:, logger:) ⇒ ActionListener
constructor
A new instance of ActionListener.
-
#start {|action| ... } ⇒ void
Start listening for actions.
-
#stop ⇒ void
Stop listening for actions.
Constructor Details
#initialize(dispatcher_client:, worker_id:, logger:) ⇒ ActionListener
Returns a new instance of ActionListener.
29 30 31 32 33 34 35 36 |
# File 'lib/hatchet/worker/action_listener.rb', line 29 def initialize(dispatcher_client:, worker_id:, logger:) @dispatcher_client = dispatcher_client @worker_id = worker_id @logger = logger @running = false @heartbeat_thread = nil @missed_heartbeats = 0 end |
Instance Method Details
#start {|action| ... } ⇒ void
This method returns an undefined value.
Start listening for actions. Blocks until stopped.
Implements retry logic with exponential backoff:
- Max 15 retries
- Reset retry counter if connection was alive > 30 seconds
- Handles GRPC::Unavailable, GRPC::DeadlineExceeded, and EOF
46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 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 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 |
# File 'lib/hatchet/worker/action_listener.rb', line 46 def start(&block) @running = true retries = 0 start_heartbeat_thread while @running && retries < MAX_RETRIES connection_start = Time.now begin @logger.info("Action listener connecting for worker #{@worker_id} (attempt #{retries + 1})") # Get the server-streaming response (ListenV2) stream = @dispatcher_client.listen(worker_id: @worker_id) # Iterate over the stream, yielding each AssignedAction stream.each do |action| break unless @running # Reset retries on successful message receipt if Time.now - connection_start > HEALTHY_CONNECTION_THRESHOLD retries = 0 connection_start = Time.now end block.call(action) end # Stream ended normally (server closed) @logger.info("Action listener stream ended normally") rescue ::GRPC::Unavailable => e @logger.warn("gRPC unavailable: #{e.message}") rescue ::GRPC::DeadlineExceeded => e @logger.warn("gRPC deadline exceeded: #{e.message}") rescue ::GRPC::Cancelled => e @logger.info("gRPC stream cancelled: #{e.message}") break unless @running rescue ::GRPC::Unknown => e @logger.warn("gRPC unknown error: #{e.message}") rescue StopIteration @logger.info("Action listener stream ended (StopIteration)") rescue Interrupt @logger.info("Action listener interrupted") @running = false break rescue StandardError => e @logger.warn("Action listener error: #{e.class}: #{e.message}") end break unless @running # Check if connection was healthy long enough to reset retries connection_duration = Time.now - connection_start if connection_duration > HEALTHY_CONNECTION_THRESHOLD retries = 0 else retries += 1 end if retries >= MAX_RETRIES @logger.error("Action listener exhausted #{MAX_RETRIES} retries. Giving up.") break end # Exponential backoff backoff = [BASE_BACKOFF_SECONDS * (2**(retries - 1)), MAX_BACKOFF_SECONDS].min @logger.info("Reconnecting in #{backoff}s (retry #{retries}/#{MAX_RETRIES})") sleep(backoff) end stop_heartbeat_thread @logger.info("Action listener stopped for worker #{@worker_id}") end |
#stop ⇒ void
This method returns an undefined value.
Stop listening for actions.
121 122 123 124 125 |
# File 'lib/hatchet/worker/action_listener.rb', line 121 def stop @running = false stop_heartbeat_thread @logger.info("Action listener stopping for worker #{@worker_id}") end |