Class: Julewire::Ractor::Destination

Inherits:
Object
  • Object
show all
Defined in:
lib/julewire/ractor/destination.rb

Overview

rubocop:disable Metrics/ClassLength -- Owns parent queue, worker lifecycle, and health.

Constant Summary collapse

DEFAULT_MAX_QUEUE =
1024
DEFAULT_REQUEST_TIMEOUT =
1

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(output:, name: :ractor, formatter: Julewire::RecordFormatter.new, encoder: Julewire::JsonEncoder.new, max_record_bytes: Core::DEFAULT_MAX_RECORD_BYTES, max_queue: DEFAULT_MAX_QUEUE, close_output: false, request_timeout: DEFAULT_REQUEST_TIMEOUT, on_drop: nil, on_failure: nil) ⇒ Destination

rubocop:disable Metrics/ParameterLists -- Destination setup mirrors core destination knobs.

Raises:

  • (ArgumentError)


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
# File 'lib/julewire/ractor/destination.rb', line 57

def initialize( # rubocop:disable Metrics/ParameterLists -- Destination setup mirrors core destination knobs.
  output:,
  name: :ractor,
  formatter: Julewire::RecordFormatter.new,
  encoder: Julewire::JsonEncoder.new,
  max_record_bytes: Core::DEFAULT_MAX_RECORD_BYTES,
  max_queue: DEFAULT_MAX_QUEUE,
  close_output: false,
  request_timeout: DEFAULT_REQUEST_TIMEOUT,
  on_drop: nil,
  on_failure: nil
)
  @name = Core::Destinations.normalize_name(name)
  @formatter = validate_callable(formatter, name: :formatter)
  @encoder = validate_callable(encoder, name: :encoder)
  Core::Destinations::Sink.validate_writeable!(output)
  Core::Validation.validate_byte_limit!(max_record_bytes, name: :max_record_bytes)
  Core::Validation.validate_non_negative_integer!(max_queue, name: :max_queue)
  Core::Validation.validate_timeout!(request_timeout, name: :request_timeout)
  raise ArgumentError, "request_timeout must be a non-negative finite Numeric" if request_timeout.nil?

  Core::Validation.validate_callable!(on_drop, name: :on_drop, allow_nil: true)
  Core::Validation.validate_callable!(on_failure, name: :on_failure, allow_nil: true)
  @output = output
  @max_record_bytes = max_record_bytes
  @max_queue = max_queue
  @close_output = close_output
  @request_timeout = request_timeout
  @on_drop = on_drop
  @on_failure = on_failure
  initialize_tracking
  start_worker
end

Instance Attribute Details

#nameObject (readonly)

Returns the value of attribute name.



55
56
57
# File 'lib/julewire/ractor/destination.rb', line 55

def name
  @name
end

Instance Method Details

#after_fork!Object



120
121
122
123
124
125
126
127
128
129
130
131
# File 'lib/julewire/ractor/destination.rb', line 120

def after_fork!
  if (@process_id - Process.pid).zero? && !close_ports(timeout: @request_timeout)
    raise Core::Error, "ractor destination worker did not stop within #{@request_timeout} seconds"
  end

  initialize_tracking
  start_worker
  self
rescue StandardError => e
  record_failure(e, phase: :after_fork)
  self
end

#close(timeout: nil) ⇒ Object



112
113
114
115
116
117
118
# File 'lib/julewire/ractor/destination.rb', line 112

def close(timeout: nil)
  timeout = lifecycle_timeout(timeout)
  @closed.set(true)
  result = request(:close, timeout: timeout, allow_closed: true)
  close_ports(timeout: timeout)
  result
end

#emit(record) ⇒ Object



91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
# File 'lib/julewire/ractor/destination.rb', line 91

def emit(record)
  increment(:received)
  return drop(:closed_dropped, record) if closed?
  return drop(:queue_full_dropped, record) unless @queue_slots.reserve

  begin
    degradation_marker = @health.degradation_marker
    @port.send({ command: :emit, degradation_marker: degradation_marker, record: record })
    increment(:queued)
  rescue StandardError => e
    release_slot
    record_failure(e, phase: :ractor_send)
    drop(:send_error, record)
  end
  nil
end

#flush(timeout: nil) ⇒ Object



108
109
110
# File 'lib/julewire/ractor/destination.rb', line 108

def flush(timeout: nil)
  @health.recover_if_successful { request(:flush, timeout: lifecycle_timeout(timeout)) }
end

#healthObject



135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
# File 'lib/julewire/ractor/destination.rb', line 135

def health
  worker = request(:health, timeout: @request_timeout)
  if worker.equal?(false)
    worker = @worker_health.get
  else
    raise TypeError, "ractor destination health must be a Hash" unless worker.instance_of?(Hash)

    Core::Integration::Protocol.validate_symbol_keys(worker)
    @worker_health.set(worker)
  end

  @health.snapshot(
    in_flight: @queue_slots.value,
    max_queue: @max_queue,
    status: status_for(worker),
    worker: worker
  )
end

#resource_identityObject



133
# File 'lib/julewire/ractor/destination.rb', line 133

def resource_identity = self