Class: RubyReactor::Executor::StepExecutor
- Inherits:
-
Object
- Object
- RubyReactor::Executor::StepExecutor
- Includes:
- AsyncStepDispatch
- Defined in:
- lib/ruby_reactor/executor/step_executor.rb
Instance Method Summary collapse
- #execute_all_steps ⇒ Object
- #execute_step(step_config) ⇒ Object
-
#initialize(context:, dependency_graph:, reactor_class:, managers:) ⇒ StepExecutor
constructor
A new instance of StepExecutor.
Constructor Details
#initialize(context:, dependency_graph:, reactor_class:, managers:) ⇒ StepExecutor
Returns a new instance of StepExecutor.
8 9 10 11 12 13 14 15 16 17 |
# File 'lib/ruby_reactor/executor/step_executor.rb', line 8 def initialize(context:, dependency_graph:, reactor_class:, managers:) @context = context @dependency_graph = dependency_graph @reactor_class = reactor_class @retry_manager = managers[:retry_manager] @result_handler = managers[:result_handler] @compensation_manager = managers[:compensation_manager] @middlewares = managers[:middlewares] || context.middlewares || Executor.middlewares_for(reactor_class) @on_step_complete = managers[:on_step_complete] end |
Instance Method Details
#execute_all_steps ⇒ Object
19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 |
# File 'lib/ruby_reactor/executor/step_executor.rb', line 19 def execute_all_steps until @dependency_graph.all_completed? || @context.finished? ready_steps = @dependency_graph.ready_steps if ready_steps.empty? raise Error::DependencyError.new( "No ready steps available but execution not complete", context: @context ) end # Execute steps sequentially ready_steps.each do |step_config| result = execute_step(step_config) # If step execution was handed off to async, return the async result return result if result.is_a?(RubyReactor::DispatchResult) # If a step returns RetryQueuedResult, we need to stop and return it return result if result.is_a?(RetryQueuedResult) # If a step returns Halt, stop the reactor cleanly (no # compensation). Must be checked BEFORE Failure / Success because # Halt is a Success subclass. return result if result.is_a?(RubyReactor::Halt) # If a step returns Failure, we need to stop execution and return it return result if result.is_a?(RubyReactor::Failure) # If a step returns InterruptResult, we need to stop execution and return it return result if result.is_a?(RubyReactor::InterruptResult) # A Skipped step (or a plain Success) continues the loop — Skipped # is a Success subclass, so this also fires the durable checkpoint # for it, same as a plain success. # # Only a continue-Success/Skipped reaches here (Async/Retry/Halt/ # Failure/Interrupt all returned above; nil is inline-async test # mode). It is the one outcome where the loop proceeds to more # steps with no other save in between — every terminal/handoff # result persists via its own path. Write a durable checkpoint so # a crash re-runs at most this one step. Ordering: side-effect -> # record result (inside execute_step) -> checkpoint here. @on_step_complete&.call if result.is_a?(RubyReactor::Success) end end # Return the final result @result_handler.final_result(@reactor_class) end |
#execute_step(step_config) ⇒ Object
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 |
# File 'lib/ruby_reactor/executor/step_executor.rb', line 70 def execute_step(step_config) if @dependency_graph.completed.include?(step_config.name) return RubyReactor.Success(@context.get_result(step_config.name)) end # Decided BEFORE argument resolution: resolving can block (or park) on # an async `result(:name)`, and when the step body is about to be # dispatched elsewhere — a `before:` hand-off, or an `async_step`'s # own worker — that wait belongs to the process that will actually run # it, not this one. (`async_reactor` still resolves here: its resolved # values are the child's INPUTS, needed at dispatch.) deferred_body = step_config.async_dispatch == :step || handoff_at?(step_config, :before) resolved_arguments = deferred_body ? {} : resolve_arguments(step_config) @middlewares.on(:start_step, step_config.name, resolved_arguments, @context) completed = false begin result = if step_config.interrupt? handle_interrupt_step(step_config) elsif step_config.async_dispatch == :step dispatch_async_step(step_config) elsif handoff_at?(step_config, :before) # `before: :x` hands off INSTEAD of running :x, leaving its # graph node incomplete for the worker to pick up. handle_background_handoff(step_config) else execute_step_with_retry(step_config, resolved_arguments) end # `after: :x` hands off once :x's result is recorded — the step really # did run here, and only what remains moves to the worker. result = handle_background_handoff(step_config) if handoff_after?(step_config, result) completed = true if result.is_a?(RubyReactor::Failure) @middlewares.on(:failed_step, step_config.name, result, @context) else @middlewares.on(:complete_step, step_config.name, result, @context) end result rescue Exception => e # rubocop:disable Lint/RescueException @middlewares.on(:failed_step, step_config.name, e, @context) unless completed raise end end |