Module: PatientHttp::Sidekiq
- Defined in:
- lib/patient_http/sidekiq.rb,
lib/patient_http/sidekiq/stats.rb,
lib/patient_http/sidekiq/web_ui.rb,
lib/patient_http/sidekiq/context.rb,
lib/patient_http/sidekiq/task_handler.rb,
lib/patient_http/sidekiq/task_monitor.rb,
lib/patient_http/sidekiq/configuration.rb,
lib/patient_http/sidekiq/request_worker.rb,
lib/patient_http/sidekiq/callback_worker.rb,
lib/patient_http/sidekiq/lifecycle_hooks.rb,
lib/patient_http/sidekiq/request_executor.rb,
lib/patient_http/sidekiq/processor_observer.rb,
lib/patient_http/sidekiq/direct_task_handler.rb,
lib/patient_http/sidekiq/task_monitor_thread.rb
Defined Under Namespace
Classes: CallbackWorker, Configuration, Context, DirectTaskHandler, LifecycleHooks, ProcessorObserver, RequestExecutor, RequestWorker, Stats, TaskHandler, TaskMonitor, TaskMonitorThread, WebUI
Constant Summary collapse
- VERSION =
File.read(File.("../../../VERSION", __FILE__)).strip
Class Attribute Summary collapse
-
.configuration ⇒ Configuration
Ensure configuration is initialized.
-
.processor ⇒ PatientHttp::Processor?
private
Returns the processor instance (internal accessor).
Class Method Summary collapse
-
.after_completion {|response| ... } ⇒ Object
Add a callback to be executed after a successful request completion.
-
.after_error {|error| ... } ⇒ Object
Add a callback to be executed after a request error.
-
.append_middleware ⇒ void
Add Sidekiq middleware for context handling.
-
.configure {|Configuration| ... } ⇒ Configuration
Configure the gem with a block.
-
.decrypt(data) ⇒ Hash
private
Decrypt data using the configured encryptor.
-
.draining? ⇒ Boolean
Check if the processor is draining (not accepting new requests but still processing in-flight ones).
-
.encrypt(data) ⇒ Hash
private
Encrypt data using the configured encryptor.
-
.execute(request, callback:, callback_args: nil, raise_error_responses: false) ⇒ String
Execute an async HTTP request.
-
.external_storage ⇒ PatientHttp::ExternalStorage
private
Get an ExternalStorage instance for storing and fetching payloads.
-
.invoke_completion_callbacks(response) ⇒ void
private
Invoke the registered completion callbacks.
-
.invoke_error_callbacks(error) ⇒ void
private
Invoke the registered error callbacks.
-
.quiet ⇒ void
Signal the processor to drain (stop accepting new requests).
-
.register_handler ⇒ void
Register Sidekiq as the request handler for processing HTTP requests.
-
.reset! ⇒ void
private
Reset all state (useful for testing).
-
.reset_configuration! ⇒ Configuration
Reset configuration to defaults (useful for testing).
-
.running? ⇒ Boolean
Check if the processor is running.
-
.start ⇒ void
Start the processor.
-
.stop(timeout: nil) ⇒ void
Stop the processor gracefully.
-
.stopped? ⇒ Boolean
Check if the processor is stopped or has not been started.
-
.stopping? ⇒ Boolean
Check if the processor is in the process of stopping.
-
.with_sidekiq_options(options) { ... } ⇒ Object
Set Sidekiq job options for HTTP requests enqueued within the block.
Class Attribute Details
.configuration ⇒ Configuration
Ensure configuration is initialized
121 122 123 |
# File 'lib/patient_http/sidekiq.rb', line 121 def configuration @configuration ||= Configuration.new end |
.processor ⇒ PatientHttp::Processor?
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Returns the processor instance (internal accessor)
405 406 407 |
# File 'lib/patient_http/sidekiq.rb', line 405 def processor @processor end |
Class Method Details
.after_completion {|response| ... } ⇒ Object
Add a callback to be executed after a successful request completion.
137 138 139 |
# File 'lib/patient_http/sidekiq.rb', line 137 def after_completion(&block) @after_completion_callbacks << block end |
.after_error {|error| ... } ⇒ Object
Add a callback to be executed after a request error.
145 146 147 |
# File 'lib/patient_http/sidekiq.rb', line 145 def after_error(&block) @after_error_callbacks << block end |
.append_middleware ⇒ void
This method returns an undefined value.
Add Sidekiq middleware for context handling. The middleware
is already added during initialization. You can call this method again to
append the middleware if needed to insert it after other middleware. If you need
further control, you can manually add the PatientHttp::Sidekiq::Context::Middleware
middleware yourself.
156 157 158 159 160 161 162 |
# File 'lib/patient_http/sidekiq.rb', line 156 def append_middleware ::Sidekiq.configure_server do |config| config.server_middleware do |chain| chain.add PatientHttp::Sidekiq::Context::Middleware end end end |
.configure {|Configuration| ... } ⇒ Configuration
Configure the gem with a block. The built configuration is also set as the
PatientHttp.default_configuration so that secrets registered at the module
level with PatientHttp.register_secret are applied to the configuration the
processor runs with, regardless of boot order.
109 110 111 112 113 114 115 116 117 |
# File 'lib/patient_http/sidekiq.rb', line 109 def configure configuration = Configuration.new yield(configuration) if block_given? @configuration = configuration @external_storage = nil register_handler PatientHttp.default_configuration = configuration configuration end |
.decrypt(data) ⇒ Hash
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Decrypt data using the configured encryptor.
215 216 217 |
# File 'lib/patient_http/sidekiq.rb', line 215 def decrypt(data) configuration.encryptor.decrypt(data) end |
.draining? ⇒ Boolean
Check if the processor is draining (not accepting new requests but still processing in-flight ones).
175 176 177 |
# File 'lib/patient_http/sidekiq.rb', line 175 def draining? !!@processor&.draining? end |
.encrypt(data) ⇒ Hash
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Encrypt data using the configured encryptor.
206 207 208 |
# File 'lib/patient_http/sidekiq.rb', line 206 def encrypt(data) configuration.encryptor.encrypt(data) end |
.execute(request, callback:, callback_args: nil, raise_error_responses: false) ⇒ String
Execute an async HTTP request.
260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 |
# File 'lib/patient_http/sidekiq.rb', line 260 def execute(request, callback:, callback_args: nil, raise_error_responses: false) PatientHttp::CallbackValidator.validate!(callback) callback_name = callback.is_a?(Class) ? callback.name : callback.to_s callback_args = PatientHttp::CallbackValidator.validate_callback_args(callback_args) request_id = SecureRandom.uuid request_json = request.as_json encrypted = encrypt(request_json) data = if external_storage.enabled? external_storage.store(encrypted, max_size: configuration.payload_store_threshold) else encrypted end = if &.any? queue = ["queue"] = .merge("patient_http_callback_queue" => queue.to_s) if queue end args = [data, callback_name, raise_error_responses, callback_args, request_id] if direct_execution?() execute_on_local_processor( request_json, args, callback_name: callback_name, raise_error_responses: raise_error_responses, callback_args: callback_args, request_id: request_id ) elsif &.any? RequestWorker.set().perform_async(*args) else RequestWorker.perform_async(*args) end request_id end |
.external_storage ⇒ PatientHttp::ExternalStorage
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Get an ExternalStorage instance for storing and fetching payloads.
197 198 199 |
# File 'lib/patient_http/sidekiq.rb', line 197 def external_storage @external_storage ||= PatientHttp::ExternalStorage.new(configuration) end |
.invoke_completion_callbacks(response) ⇒ void
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
This method returns an undefined value.
Invoke the registered completion callbacks
384 385 386 387 388 |
# File 'lib/patient_http/sidekiq.rb', line 384 def invoke_completion_callbacks(response) @after_completion_callbacks.each do |callback| callback.call(response) end end |
.invoke_error_callbacks(error) ⇒ void
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
This method returns an undefined value.
Invoke the registered error callbacks
395 396 397 398 399 |
# File 'lib/patient_http/sidekiq.rb', line 395 def invoke_error_callbacks(error) @after_error_callbacks.each do |callback| callback.call(error) end end |
.quiet ⇒ void
This method returns an undefined value.
Signal the processor to drain (stop accepting new requests)
338 339 340 341 342 343 344 |
# File 'lib/patient_http/sidekiq.rb', line 338 def quiet @lifecycle_mutex.synchronize do return unless running? @processor.drain end end |
.register_handler ⇒ void
This method returns an undefined value.
Register Sidekiq as the request handler for processing HTTP requests. This is called automatically when the processor starts or you call PatientHttp::Sidekiq.configure.
304 305 306 307 308 309 310 311 312 313 314 315 |
# File 'lib/patient_http/sidekiq.rb', line 304 def register_handler @request_handler ||= lambda do |request:, callback:, raise_error_responses:, callback_args:| execute( request, callback: callback, raise_error_responses: raise_error_responses, callback_args: callback_args ) end PatientHttp.register_handler(@request_handler) end |
.reset! ⇒ void
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
This method returns an undefined value.
Reset all state (useful for testing)
367 368 369 370 371 372 373 374 375 376 377 |
# File 'lib/patient_http/sidekiq.rb', line 367 def reset! @lifecycle_mutex.synchronize do @processor&.stop(timeout: 0) @processor = nil end @configuration = nil @external_storage = nil @after_completion_callbacks = [] @after_error_callbacks = [] PatientHttp.unregister_handler(@request_handler) if @request_handler end |
.reset_configuration! ⇒ Configuration
Reset configuration to defaults (useful for testing)
127 128 129 130 131 |
# File 'lib/patient_http/sidekiq.rb', line 127 def reset_configuration! @configuration = nil @external_storage = nil configuration end |
.running? ⇒ Boolean
Check if the processor is running.
167 168 169 |
# File 'lib/patient_http/sidekiq.rb', line 167 def running? !!@processor&.running? end |
.start ⇒ void
This method returns an undefined value.
Start the processor
320 321 322 323 324 325 326 327 328 329 330 331 332 333 |
# File 'lib/patient_http/sidekiq.rb', line 320 def start @lifecycle_mutex.synchronize do return if @processor && !@processor.stopped? @processor = PatientHttp::Processor.new(configuration) @processor.observe(ProcessorObserver.new(@processor)) configuration.observers.each do |observer| @processor.observe(observer) end @processor.start end register_handler end |
.stop(timeout: nil) ⇒ void
This method returns an undefined value.
Stop the processor gracefully
350 351 352 353 354 355 356 357 358 359 360 361 |
# File 'lib/patient_http/sidekiq.rb', line 350 def stop(timeout: nil) if @request_handler PatientHttp.unregister_handler(@request_handler) end @lifecycle_mutex.synchronize do return unless @processor @processor.stop(timeout: timeout) @processor = nil end end |
.stopped? ⇒ Boolean
Check if the processor is stopped or has not been started.
189 190 191 |
# File 'lib/patient_http/sidekiq.rb', line 189 def stopped? @processor.nil? || @processor.stopped? end |
.stopping? ⇒ Boolean
Check if the processor is in the process of stopping.
182 183 184 |
# File 'lib/patient_http/sidekiq.rb', line 182 def stopping? !!@processor&.stopping? end |
.with_sidekiq_options(options) { ... } ⇒ Object
Set Sidekiq job options for HTTP requests enqueued within the block.
Options are applied with Sidekiq's set method (queue, retry, etc.).
Nested calls merge options with the innermost values taking precedence.
If the options include a queue, the callback job for the request is
enqueued on that queue as well. Options only apply to requests enqueued
in the same fiber as the block. Requests made in the block always go
through the Sidekiq queue, even when direct execution is enabled, so
that Sidekiq applies the options. This method has no effect
when jobs run inline with Sidekiq::Testing.inline!.
232 233 234 235 236 237 238 239 240 241 242 243 244 245 |
# File 'lib/patient_http/sidekiq.rb', line 232 def () unless .is_a?(Hash) raise ArgumentError.new("options must be a Hash, got: #{.class}") end raise ArgumentError.new("with_sidekiq_options requires a block") unless block_given? previous = Thread.current[:patient_http_sidekiq_options] begin Thread.current[:patient_http_sidekiq_options] = (previous || {}).merge(.transform_keys(&:to_s)) yield ensure Thread.current[:patient_http_sidekiq_options] = previous end end |