Class: Hatchet::WorkerRuntime::ActionListener

Inherits:
Object
  • Object
show all
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.

Examples:

listener = ActionListener.new(dispatcher_client: client, worker_id: id, logger: logger)
listener.start do |action|
  runner.execute(action)
end

Constant Summary collapse

MAX_RETRIES =

Returns:

  • (Integer)
15
BASE_BACKOFF_SECONDS =

Returns:

  • (Integer)
1
MAX_BACKOFF_SECONDS =

Returns:

  • (Integer)
30
HEALTHY_CONNECTION_THRESHOLD =

seconds

Returns:

  • (Integer)
30
HEARTBEAT_INTERVAL =

seconds

Returns:

  • (Integer)
4
MAX_MISSED_HEARTBEATS =

Returns:

  • (Integer)
3

Instance Method Summary collapse

Constructor Details

#initialize(dispatcher_client:, worker_id:, logger:) ⇒ ActionListener

Returns a new instance of ActionListener.

Parameters:



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

Yields:

  • (action)

    Called for each action received from the dispatcher

Yield Parameters:

  • action (Object)

Yield Returns:

  • (void)


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

#stopvoid

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