Class: Julewire::Ractor::Destination
- Inherits:
-
Object
- Object
- Julewire::Ractor::Destination
- 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
-
#name ⇒ Object
readonly
Returns the value of attribute name.
Instance Method Summary collapse
- #after_fork! ⇒ Object
- #close(timeout: nil) ⇒ Object
- #emit(record) ⇒ Object
- #flush(timeout: nil) ⇒ Object
- #health ⇒ Object
-
#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
constructor
rubocop:disable Metrics/ParameterLists -- Destination setup mirrors core destination knobs.
- #resource_identity ⇒ Object
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.
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
#name ⇒ Object (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 |
#health ⇒ Object
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_identity ⇒ Object
133 |
# File 'lib/julewire/ractor/destination.rb', line 133 def resource_identity = self |